我在使用官方 Confluent Kafka Python API 时遇到错误:
我订阅:
kafka_consumer.subscribe(topics=["my-avro-topic"], on_assign=on_assign_callback, on_revoke=on_revoke_callback)
Run Code Online (Sandbox Code Playgroud)
使用回调:
def on_assign_callback(consumer, topic_partitions):
for topic_partition in topic_partitions:
print("without position. topic={}. partition={}. offset={}. error={}".format(topic_partition.topic, topic_partition.partition,
topic_partition.offset, topic_partition.error))
topic_partitions_with_offsets = consumer.position(topic_partitions)
print("assigned to {}->{} partitions".format(len(topic_partitions), len(topic_partitions_with_offsets)))
for topic_partition in topic_partitions_with_offsets:
print("with position. topic={}. partition={}. offset={}. error={}".format(topic_partition.topic, topic_partition.partition,
topic_partition.offset, topic_partition.error))
Run Code Online (Sandbox Code Playgroud)
产生控制台输出:
without position. topic=my-avro-topic. partition=0. offset=-1001. error=None
assigned to 1->1 partitions
with position. topic=my-avro-topic. partition=0. offset=-1001. error=KafkaError{code=_UNKNOWN_PARTITION,val=-190,str="(null)"}
Run Code Online (Sandbox Code Playgroud)
有人可以解释一下吗?为什么我会收到未知分区的回调通知?类似的代码使用 Java API 可以完美运行。
我创建了很多没有密钥的主题,如何修改它们并添加正确的主题?
我需要为一些希望他们正确阅读主题的连接器更改此设置
我个人使用 ksql 但我没有找到任何方法来做到这一点
我已经阅读了他们的博客并理解了他们的例子。 https://www.confluent.io/blog/apache-kafka-to-amazon-s3-exactly-once/
但我正在努力解决我所拥有的这种情况。我目前的配置是:
"flush.size": "50",
"rotate.interval.ms": "-1",
"rotate.schedule.interval.ms": "300000",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "3600000",
"path.format": "YYYY/MM/dd/HH",
"timestamp.extractor": "Wallclock"
Run Code Online (Sandbox Code Playgroud)
根据我对配置的了解。连接器将50在300000ms(5 分钟)后提交记录文件或文件,以先到者为准。如果连接器将文件上传到 s3 但未能提交到 Kafka,由于我设置了轮换计划间隔,Kafka 如何重新上传将覆盖 s3 文件的相同记录?这不会导致 s3 重复吗?
amazon-s3 apache-kafka apache-kafka-connect confluent-platform
就投资最大价值而言,在托管端到端 Kafka 事件源方面,AWS MSK 与 Confluence 相比如何?
用于比较的主要标准是: