我有一个部署,我们使用 kafka 从服务发送消息。但是我们需要在所有地区都拥有 Kafka 大师。因此,一旦消息在 1 个数据中心推送,它应该在其他数据中心同步。当它再次在其他数据中心完成时,它应该同步回来。Mirror Maker 可以提供从 1 到其他的同步,但我如何实现双向同步?
我有一个有 40 个分区的主题。设置是这样的:
def on_assign (c,ps):
for p in ps:
p.offset=0
print ps
c.assign(ps)
conf = {'bootstrap.servers': 'localhost:9092'
'enable.auto.commit' : False,
'group.id' : 'confluent_consumer',
'default.topic.config': {'auto.offset.reset': 'earliest'}
}
consumer = Consumer(**conf)
consumer.subscribe(['topic.source'], on_assign=on_assign)
msg = consumer.poll(timeout=100000)
print "Topic is %s: | Partition is %d: | Offset is : %d | key is :%s " % (msg.topic(), msg.partition(), msg.offset(), msg.key())
Run Code Online (Sandbox Code Playgroud)
我想从偏移量 0 读取主题的所有分区topic.source。但我没有看到所有分区都会发生这种情况。对于某些分区,它从特定的偏移量读取,我假设这是提交的偏移量,group.id每次更改也没有帮助。如何从头开始读取该主题的所有分区,而不考虑提交的偏移量?
我打印ps出来on_assign(),它为所有 40 个分区打印了这样的内容:
[TopicPartition{topic=topic.source,partition=0,offset=0,error=None},TopicPartition{topic=topic.source,partition=1,offset=0,error=None}....] and so on
Run Code Online (Sandbox Code Playgroud) 尝试使用以下代码构建一个kafka消费者
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
// set up consumer
final Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, CLUSTER.bootstrapServers());
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "consumer-tutorial");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
io.confluent.kafka.serializers.KafkaAvroSerializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
io.confluent.kafka.serializers.KafkaAvroSerializer.class);
// transactional API
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
// consumer --from-beginning
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
consumerProps.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "10000");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
consumerProps.put("zookeeper.connect", CLUSTER.zookeeperConnect());
consumerProps.put("schema.registry.url", CLUSTER.schemaRegistryUrl());
final KafkaConsumer<GenericRecord, GenericRecord> consumer = new KafkaConsumer<GenericRecord, GenericRecord>(consumerProps);
consumer.subscribe(Collections.singletonList(inputTopic));
Run Code Online (Sandbox Code Playgroud)
但因错误而失败
org.apache.kafka.common.KafkaException: Failed to construct kafka consumer
at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:765)
at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:633)
at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:615)
at com.telefonica.app.test_consumer.KafkaETLConsumerTest.testRunConsumer(KafkaETLConsumerTest.java:192)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) …Run Code Online (Sandbox Code Playgroud) 是否可以在 Windows 上运行完整的 Confluence 平台?如果不是,运行 Confluence 平台的最佳方式是什么?
谢谢
我创建了最简单的 kafka 接收器连接器配置,我使用的是 confluence 4.1.0:
{
"connector.class":
"io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"type.name": "test-type",
"tasks.max": "1",
"topics": "dialogs",
"name": "elasticsearch-sink",
"key.ignore": "true",
"connection.url": "http://localhost:9200",
"schema.ignore": "true"
}
Run Code Online (Sandbox Code Playgroud)
在主题中,我将消息保存为JSON
{ "topics": "resd"}
Run Code Online (Sandbox Code Playgroud)
但在结果中我得到一个错误:
引起:org.apache.kafka.common.errors.SerializationException:反序列化 id -1 的 Avro 消息时出错 引起:org.apache.kafka.common.errors.SerializationException:未知的魔术字节!
我设置了 kafka 连接 s3 接收器,持续时间设置为 1 小时,并且我设置了一个相当大的刷新计数,比如 10,000。现在如果kafka通道中的消息不多,s3 sink会尝试将它们缓存在内存中,等待它累积到flush计数,然后将它们一起上传并将偏移量提交给自己的消费者组。
但是想想这种情况。如果在频道里,我只发送5000条消息。然后没有 s3 水槽冲洗。然后时间长了,这5000条消息最终会因为保留时间的原因从kafka中驱逐出去。但是这些消息仍然在 s3 sink 的内存中,而不是在 s3 中。这是非常危险的,例如,如果我们重新启动 s3 sink 或运行 s3 sink 的机器就崩溃了。然后我们丢失了那 5,000 条消息。我们无法从 kafka 中再次找到它们,因为它已经被删除了。
这会发生在 s3 sink 上吗?或者有一些设置会强制它在一段时间后刷新?
amazon-s3 apache-kafka apache-kafka-connect confluent-platform
我的主要问题:为什么模式注册表会崩溃?
外围问题:如果我为每个 zookeeper/kafka/schema-registry 配置了一个服务器,为什么每个 pod 都启动?其他一切看起来基本正确吗?
? helm repo update
<snip>
? helm install --values values.yaml --name my-confluent-oss confluentinc/cp-helm-charts
<snip>
? helm list
NAME REVISION UPDATED STATUS CHART APP VERSION NAMESPACE
my-confluent-oss 1 Sat Oct 20 19:09:08 2018 DEPLOYED cp-helm-charts-0.1.0 1.0 default
? kubectl get pods
NAME READY STATUS RESTARTS AGE
my-confluent-oss-cp-kafka-0 2/2 Running 0 20m
my-confluent-oss-cp-schema-registry-59d8877584-c2jc7 1/2 CrashLoopBackOff 7 20m
my-confluent-oss-cp-zookeeper-0 2/2 Running 0 20m
Run Code Online (Sandbox Code Playgroud)
我values.yaml的如下。我已经用helm install --debug --dry-run. 我只是禁用持久性,设置单个服务器(这是在 VM 中运行的开发设置),并暂时禁用额外服务,直到我获得基础工作:
cp-kafka:
brokers: 1 …Run Code Online (Sandbox Code Playgroud) apache-kafka kubernetes kubernetes-helm confluent-schema-registry confluent-platform
我正在运行 kafka-connect 分布式设置。
我正在使用单机/进程设置(仍处于分布式模式)进行测试,效果很好,现在我正在使用 3 个节点(和 3 个连接进程),日志不包含错误,但是当我提交 s3-connector 时通过rest-api请求,它返回:{"error_code":409,"message":"Cannot complete request because of a conflicting operation (e.g. worker rebalance)"}。
当我停止其中一个节点上的 kafka-connect 进程时,我实际上可以提交作业并且一切正常。
我的集群中有 3 个代理,主题的分区号是 32。
这是我尝试启动的连接器:
{
"name": "s3-sink-new-2",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "32",
"topics": "rawEventsWithoutAttribution5",
"s3.region": "us-east-1",
"s3.bucket.name": "dy-raw-collection",
"s3.part.size": "64000000",
"flush.size": "10000",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.avro.AvroFormat",
"schema.generator.class": "io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "60000",
"path.format": "\'year\'=YYYY/\'month\'=MM/\'day\'=dd/\'hour\'=HH",
"locale": "US",
"timezone": "GMT",
"timestamp.extractor": "RecordField",
"timestamp.field": "procTimestamp",
"name": "s3-sink-new-2"
}
}
Run Code Online (Sandbox Code Playgroud)
日志中没有任何内容表明有问题,我真的迷失在这里。
我想从融合的云主题读取数据,然后写入另一个主题。
在本地主机上,我没有遇到任何重大问题。但是confluent cloud的schema registry需要传递一些我不知道如何输入的身份验证数据:
basic.auth.credentials.source=USER_INFO
schema.registry.basic.auth.user.info=:
schema.registry.url= https://xxxxxxxxxx.confluent.cloudBlockquote
以下是当前代码:
import com.databricks.spark.avro.SchemaConverters
import io.confluent.kafka.schemaregistry.client.{CachedSchemaRegistryClient, SchemaRegistryClient}
import io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer
import org.apache.avro.Schema
import org.apache.avro.generic.GenericRecord
import org.apache.spark.sql.SparkSession
object AvroConsumer {
private val topic = "transactions"
private val kafkaUrl = "http://localhost:9092"
private val schemaRegistryUrl = "http://localhost:8081"
private val schemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 128)
private val kafkaAvroDeserializer = new AvroDeserializer(schemaRegistryClient)
private val avroSchema = schemaRegistryClient.getLatestSchemaMetadata(topic + "-value").getSchema
private var sparkSchema = SchemaConverters.toSqlType(new Schema.Parser().parse(avroSchema))
def main(args: Array[String]): Unit = {
val spark = SparkSession
.builder
.appName("ConfluentConsumer")
.master("local[*]")
.getOrCreate() …Run Code Online (Sandbox Code Playgroud) apache-spark databricks confluent-schema-registry spark-structured-streaming confluent-platform
要开发我的 Kafka 连接器,我需要添加一个 connect-API 依赖项。
我应该使用哪一种?
例如 mongodb 连接器使用来自maven central 的connect-api
但是来自开发指南的链接转到https://packages.confluent.io/maven/org/apache/kafka/connect-api/5.5.0-ccs/旁边5.5.0-ccs还有5.5.0-ce版本。
所以,此时最后的版本是:
所有三个变体之间有什么区别?
我应该使用哪一种?
java maven apache-kafka apache-kafka-connect confluent-platform
apache-kafka ×9
amazon-s3 ×1
apache-spark ×1
databricks ×1
java ×1
kubernetes ×1
maven ×1
python ×1
windows ×1