You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Copy file name to clipboardexpand all lines: floodplain-dsl/src/main/kotlin/io/floodplain/kotlindsl/Stream.kt
+2-25
Original file line number
Diff line number
Diff line change
@@ -304,40 +304,17 @@ class Stream(override val topologyContext: TopologyContext, val topologyConstruc
304
304
connectorClientConfigOverridePolicy
305
305
)
306
306
307
-
val statusBackingStore:StatusBackingStore=KafkaStatusBackingStore(time, worker.getInternalValueConverter())
307
+
val statusBackingStore:StatusBackingStore=KafkaStatusBackingStore(time, worker.internalValueConverter)
308
308
statusBackingStore.configure(config)
309
309
310
-
val configBackingStore =KafkaConfigBackingStore(worker.getInternalValueConverter(), config, worker.configTransformer())
310
+
val configBackingStore =KafkaConfigBackingStore(worker.internalValueConverter, config, worker.configTransformer())
311
311
val herder =DistributedHerder(config, time, worker, kafkaClusterId, statusBackingStore, configBackingStore, advertisedUrl.toString(), connectorClientConfigOverridePolicy)
312
312
313
-
// val herder: Herder = StandaloneHerder(worker, kafkaClusterId, connectorClientConfigOverridePolicy)
314
313
val connect =Connect(herder, rest)
315
314
316
315
connect.start()
317
316
logger.info("Connect started!!")
318
317
return herder
319
-
// try {
320
-
// // for (connectorPropsFile in Arrays.copyOfRange(args, 1, args.length)) {
321
-
// // val connectorProps: Map<String, String> = Utils.propsToStringMap(Utils.loadProps(connectorPropsFile))
322
-
// // val cb: FutureCallback<Herder.Created<ConnectorInfo>> =
0 commit comments