This reverts commit 8281cbf762.
This commit is contained in:
@@ -111,8 +111,6 @@ public class BulkInsertWriteFunction<I>
|
||||
|
||||
@Override
|
||||
public void open(Configuration parameters) throws IOException {
|
||||
// always use the user classloader
|
||||
Thread.currentThread().setContextClassLoader(getRuntimeContext().getUserCodeClassLoader());
|
||||
this.taskID = getRuntimeContext().getIndexOfThisSubtask();
|
||||
this.metaClient = StreamerUtil.createMetaClient(this.config);
|
||||
this.writeClient = StreamerUtil.createWriteClient(this.config, getRuntimeContext());
|
||||
|
||||
@@ -125,8 +125,6 @@ public abstract class AbstractStreamWriteFunction<I>
|
||||
|
||||
@Override
|
||||
public void initializeState(FunctionInitializationContext context) throws Exception {
|
||||
// always use the user classloader
|
||||
Thread.currentThread().setContextClassLoader(getRuntimeContext().getUserCodeClassLoader());
|
||||
this.taskID = getRuntimeContext().getIndexOfThisSubtask();
|
||||
this.metaClient = StreamerUtil.createMetaClient(this.config);
|
||||
this.writeClient = StreamerUtil.createWriteClient(this.config, getRuntimeContext());
|
||||
|
||||
@@ -75,8 +75,6 @@ public class CompactFunction extends ProcessFunction<CompactionPlanEvent, Compac
|
||||
|
||||
@Override
|
||||
public void open(Configuration parameters) throws Exception {
|
||||
// always use the user classloader
|
||||
Thread.currentThread().setContextClassLoader(getRuntimeContext().getUserCodeClassLoader());
|
||||
this.taskID = getRuntimeContext().getIndexOfThisSubtask();
|
||||
this.writeClient = StreamerUtil.createWriteClient(conf, getRuntimeContext());
|
||||
if (this.asyncCompaction) {
|
||||
|
||||
Reference in New Issue
Block a user