几天以来,我一直在玩融合版本的 kafka,以更好地了解平台。对于发送到一个主题的某些格式错误的 avro 消息,我收到了一些序列化异常。让我用事实来解释这个问题:
<kafka.new.version>0.10.2.0-cp1</kafka.new.version>
<confluent.version>3.2.0</confluent.version>
<avro.version>1.7.7</avro.version>
Run Code Online (Sandbox Code Playgroud)
意图:非常简单,Producer 发送 Avro 记录,Consumer 应该毫无问题地消费所有记录,(它可以使所有消息与架构注册表中的架构不兼容。)用法:
Producer ->
Key -> StringSerializer
Value -> KafkaAvroSerializer
Consumer ->
Key -> StringDeserializer
Value -> KafkaAvroDeserializer
Run Code Online (Sandbox Code Playgroud)
其他消费者属性(仅供参考):
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "somehost:9092");
properties.put(ConsumerConfig.GROUP_ID_CONFIG, "myconsumer-4");
properties.put(ConsumerConfig.CLIENT_ID_CONFIG, "someclient-4");
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringDeserializer.class);
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, io.confluent.kafka.serializers.KafkaAvroDeserializer.class);
properties.put(AUTO_OFFSET_RESET_CONFIG, "earliest");
properties.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true);
properties.put("schema.registry.url", "schemaregistryhost:8081");
Run Code Online (Sandbox Code Playgroud)
我能够毫无问题地使用消息,直到其他一些生产者错误地向该主题发送了一条消息并修改了架构注册表中的最新架构。(我们在架构注册表中启用了一个选项,因此您可以向主题发送任何消息,架构注册表每次都会创建一个新版本的架构,如果关闭,我们也可以关闭。)
现在,由于这一个坏消息,则poll()是序列化问题失败。它确实给了我失败的偏移量,我可以通过使用 seek() 传递偏移量,但这听起来不太好。我还尝试使用最大轮询记录为 10 并将 poll() 超时设置为非常小,以便我可以通过捕获异常来忽略最多 10 条记录,但由于某种原因,max-records 不起作用并且代码立即失败并出现序列化错误,即使我从开始和坏消息在 240 偏移处。
properties.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "10");
Run Code Online (Sandbox Code Playgroud)
另一个简单的解决方案是在我的应用程序中使用 ByteArrayDeserializer 并使用 KafkaAvroDecoder,我可以处理反序列化问题。
我相信我缺少某些东西或做错了。也添加例外:
Exception in thread "main" org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition topic.ongo.test3.user14-0 …Run Code Online (Sandbox Code Playgroud) 我想完全卸载融合。我按照他们网站上的说明安装了它。有三个简单的步骤:
$ wget -qO - https://packages.confluent.io/deb/4.0/archive.key | sudo apt-key add -
$ sudo add-apt-repository "deb [arch=amd64] https://packages.confluent.io/deb/4.0 stable main"
$ sudo apt-get update && sudo apt-get install confluent-platform-oss-2.11
Run Code Online (Sandbox Code Playgroud)
现在我该如何删除/卸载它。我找不到任何与之相关的东西。
我有一个现有的 Kafka 集群。我想安装 Kafka REST 代理:
https://github.com/confluentinc/kafka-rest
如果我安装 confluent 会随 Kafka 一起出现吗?我担心如果我仍然在我的主 Kafka 节点上使用 confluent 会覆盖我的所有设置并弄乱我的 Kafka 集群。
当您有一个现有的 Kafka 集群时,如何安装 Kafka REST?这在他们的网站上没有明确说明。我有 CentOS 并打算尝试:
sudo yum install confluent-platform-oss-2.11
Run Code Online (Sandbox Code Playgroud)
任何帮助都会很棒......
我为本地运行设置了 Kafka。我已经用 Java 编写了示例生产者和消费者,并通过启动服务器和动物园管理员从本地运行。
我想使用oracle作为生产者,这需要编写配置文件(已经编写),confluent shell script才能在Unix上运行它。
有什么办法可以confluent在 Windows上运行,我confluent在安装程序中找不到批处理文件?
另外,有没有办法在不使用confluent脚本的情况下以生产者身份运行 Oracle ?
java apache-kafka kafka-producer-api apache-kafka-connect confluent-platform
有没有办法在 Kafka Connect 启动时自动加载(多个)Kafka Connect 连接器(例如在 Confluence Platform 中)?
到目前为止我发现了什么:
Confluence Docs 声明使用bin/connect-standalone
独立模式的命令以及工作线程和每个连接器的属性文件。
对于分布式模式,您必须通过 REST API 运行连接器。
https://docs.confluence.io/current/connect/userguide.html#standalone-mode,https://docs.confluence.io/current/connect/managing/configuring.html#standalone-example
是否有另一种方法,例如包含应在“connect-[standalone|distributed].properties”文件中运行的所有连接器(类似于在 ksql-server.properties 中提供 KSQL 查询文件),以便它们自动加载到Kafka Connect 的启动(例如在 Confluence 平台中)?
或者,即使在生产环境中,连接器是否也如上所述“手动”加载?
Confluence 控制中心未启动。
我执行了以下命令来启动 Confluence 平台
模式注册表启动(终端 3)
然后我尝试启动控制中心(4 号航站楼)。但出现错误。它没有开始。
[2019-10-30 08:58:36,331] INFO [main] unable to get command store (io.confluent.command.CommandStore)
[2019-10-30 08:58:37,331] INFO [main] unable to get command store (io.confluent.command.CommandStore)
[2019-10-30 08:58:37,331] WARN [main] unable to start with allowance=300000 (io.confluent.command.CommandStore)
[2019-10-30 08:58:37,332] ERROR [main] failed to start topology (io.confluent.controlcenter.ControlCenter)
java.util.concurrent.TimeoutException
at io.confluent.command.CommandStore.start(CommandStore.java:108)
at io.confluent.controlcenter.ControlCenter.main(ControlCenter.java:124)
Run Code Online (Sandbox Code Playgroud) 我们可以在视图或物化视图上设置 CDC 跟踪吗?我们使用的是 Confluence 平台,来源是 SQL Server,CDC 工具是 Debezium。
Apache Kafka随着 Kafka 2.4 的发布引入了Mirrormaker2 (MM2)。MM2明显优于MM1。
我知道从架构的角度来看,MM1 过去使用生产者和消费者 API 工作,而 MM2 使用连接 API。我相信MM2的设计灵感来自于Confluence Replicator。Confluence Replicator 与 Confluence 工具完美集成。但除此之外,MM2 和 Confluence Replicator 之间有什么区别?
我使用 conluent jdbc-sink 将数据从 kafka 加载到 oracle。
但我用数据来写我的价值模式。
我不想用数据编写模式,如何在 kafka 主题上编写模式,然后我只想从客户端发送数据?
提前致谢
json数据
{
"schema": {
"type": "struct",
"fields": [
{
"field": 'ID',
"type": "int32",
"optional": False
},
{
"field": 'PRODUCT',
"type": "string",
"optional": True
},
{
"field": 'QUANTITY',
"type": "int32",
"optional": True
},
{
"field": 'PRICE',
"type": "int32",
"optional": True
}
],
"optional": True,
"name": "myrecord"
},
"payload": {
"ID": 1071,
"PRODUCT": 'ersin',
"QUANTITIY": 1071,
"PRICE": 1453
}
Run Code Online (Sandbox Code Playgroud)
蟒蛇代码:
producer.send(topic, key=b'1071'
, value=json.dumps(v, default=json_util.default).encode('utf-8'))
Run Code Online (Sandbox Code Playgroud)
我该如何解决这个问题?
提前致谢
我在 Kafka 主题中存储了以下批处理格式:
Data generated --> B2E, T3, T2, T1, B2S | B1E, T3, T2, T1, B1S --> Data Consumed
Run Code Online (Sandbox Code Playgroud)
这里BS表示批次开始,BE表示批次结束,t1,t2,t3是三个不同的数据。
这是扫描仪:
在将批次结束记录推送到B1E主题之前,我们经过一些验证后得知该B1批次无效。
在这种情况下,B1批次范围内BS的所有数据都BE应该转到特定主题。
因此,在上面的示例中,b1批次应转到主题 T1,b2批次应转到主题 T2。
我怎样才能使用卡夫卡做到这一点?