KafkaCluster

class KafkaCluster(val kafkaConfigMap: KafkaConfigMap, pollDuration: Duration, val schemaRegistryUrl: String? = null, val transactionalIdPrefix: String? = null, groupId: String = "xtdb", coroutineContext: CoroutineContext = Dispatchers.Default) : Remote

Constructors

Link copied to clipboard
constructor(kafkaConfigMap: KafkaConfigMap, pollDuration: Duration, schemaRegistryUrl: String? = null, transactionalIdPrefix: String? = null, groupId: String = "xtdb", coroutineContext: CoroutineContext = Dispatchers.Default)

Types

Link copied to clipboard
@Serializable
@SerialName(value = "!Kafka")
data class ClusterFactory @JvmOverloads constructor(val bootstrapServers: String, var pollDuration: Duration = Duration.ofSeconds(1), var propertiesMap: Map<String, String> = emptyMap(), var propertiesFile: Path? = null, var schemaRegistryUrl: String? = null, var transactionalIdPrefix: String? = null, var groupId: String = "xtdb", var coroutineContext: CoroutineContext = Dispatchers.Default) : Remote.Factory<KafkaCluster>
Link copied to clipboard
@Serializable
@SerialName(value = "!Kafka")
data class LogFactory @JvmOverloads constructor(val cluster: RemoteAlias, val topic: String, var replicaCluster: RemoteAlias = cluster, var replicaTopic: String = "-replica", var autoCreateTopic: Boolean = true, var epoch: Int = 0) : Log.Factory

Properties

Link copied to clipboard
val kafkaConfigMap: KafkaConfigMap
Link copied to clipboard
val producer: KafkaProducer<Unit?, ByteArray?>
Link copied to clipboard
Link copied to clipboard
val scope: CoroutineScope
Link copied to clipboard

Functions

Link copied to clipboard
open override fun close()