标签: confluent-platform

Kafka:轮询期间的反序列化问题

几天以来,我一直在玩融合版本的 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)

apache-kafka kafka-consumer-api confluent-platform

5
推荐指数
1
解决办法
5972
查看次数

在 linux 上完全卸载 confluent

我想完全卸载融合。我按照他们网站上的说明安装了它。有三个简单的步骤:

$ 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)

现在我该如何删除/卸载它。我找不到任何与之相关的东西。

apt apache-kafka confluent-platform

5
推荐指数
1
解决办法
2655
查看次数

Kafka Confluent REST API:包括 Kafka?

我有一个现有的 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)

任何帮助都会很棒......

apache-kafka kafka-rest confluent-platform

5
推荐指数
1
解决办法
353
查看次数

Kafka:在 Windows 环境中运行 Confluent

我为本地运行设置了 Kafka。我已经用 Java 编写了示例生产者和消费者,并通过启动服务器和动物园管理员从本地运行。
我想使用oracle作为生产者,这需要编写配置文件(已经编写),confluent shell script才能在Unix上运行它。

有什么办法可以confluent在 Windows上运行,我confluent在安装程序中找不到批处理文件?

另外,有没有办法在不使用confluent脚本的情况下以生产者身份运行 Oracle ?

java apache-kafka kafka-producer-api apache-kafka-connect confluent-platform

5
推荐指数
1
解决办法
1万
查看次数

带有 kafka-avro-console-consumer 的未知魔法字节

我一直在尝试将 kafka-avro-console-consumer 从 Confluent 连接到我们遗留的 Kafka 集群,该集群是在没有 Confluent Schema Registry 的情况下部署的。我使用以下属性显式提供了架构:

kafka-console-consumer --bootstrap-server kafka02.internal:9092 \
    --topic test \
    --from-beginning \
    --property key.schema='{"type":"long"}' \
    --property value.schema='{"type":"long"}'
Run Code Online (Sandbox Code Playgroud)

但我收到“未知的魔法字节!” 错误org.apache.kafka.common.errors.SerializationException

是否可以使用 Confluent kafka-avro-console-consumer 消费来自 Kafka 的 Avro 消息,这些消息未使用 Confluent 的 AvroSerializer 和 Schema Registry 序列化?

avro apache-kafka confluent-schema-registry confluent-platform

5
推荐指数
1
解决办法
1万
查看次数

.NET Confluent Kafka 消费者内存泄漏

当使用 Confluent Kafka .NET 库从 Kafka 消费时,我们看到了巨大的内存泄漏。

我注意到的一件事是代码在没有 using 语句的情况下使用:

while (true)
{
    if (_consumer.Consume(out Message<string, string> message, TimeSpan.FromMilliseconds(100)))
    {
        OnMessage(message);
    }
}
Run Code Online (Sandbox Code Playgroud)

但是,该 while 循环在应用程序的整个生命周期内运行,因此 .Dispose() 永远不会在消费时被调用。不会创建其他消费者实例。

由于库背后的代码是 C 语言,如果我们调用 GC.Collect() 会清除库创建的对象,还是垃圾收集器无法控制的非托管代码?

其他任何可能导致泄漏的事情,是否需要在某些时期或类似的时间调用 consumer.Close() ?

.net c# apache-kafka confluent-platform

5
推荐指数
0
解决办法
705
查看次数

Kafka - C# - confluent-kafka-dotnet - 消息超时

我们在 Linux 机器上有一个简单的 Kafka 1.1.0 独立模式部署。在 server.properties 中,我们修改了:

listeners = PLAINTEXT://10.0.5.66:9092
Run Code Online (Sandbox Code Playgroud)

advertised.listeners被注释掉,因此它将回退到listeners属性中找到的默认值。

我们正在使用 .NET (C#) 生产者,它通过 confluent-kafka-dotnet (0.11.4) 推送消息。有时消息会传输到 Kafka,有时我们会在生产者端收到“消息超时”错误。对于可能导致此问题的原因,我们没有任何想法。它不时发生。如果一条消息失败,另一条消息通常会在第一条消息通过后几秒钟出现。

另一个迹可从我们不时在服务器上的日志,卡夫卡看到以下消息:WARN: Attempting to send a response via a channel for which there is no open connection <IP:PORT>。此消息有时包含生产者的 IP 地址和端口。

知道可能有什么问题吗?

c# timeout apache-kafka confluent-platform

5
推荐指数
0
解决办法
1401
查看次数

Confluent 平台 Kafka Connect 因退出 137 崩溃

在 Mac 上,我提取了最新的 docker 镜像。当我运行堆栈时,一切似乎都没问题,但“连接”在退出 137 时崩溃了。

当我查看指挥中心时,集群健康状况似乎很好。这有什么影响?如何纠正问题?

在此处输入图片说明

感谢任何帮助。

谢谢 !

apache-kafka apache-kafka-connect confluent-platform

5
推荐指数
1
解决办法
1833
查看次数

如何在 docker 上使用 Confluent CLI

我已经在https://docs.confluent.io/current/quickstart/ce-docker-quickstart.html的帮助下使用 docker 在我的 Windows 10 上启动了 Confluent Platform 。现在我想尝试使用 Confluent CLI。但是我没有看到任何关于如何在 docker 上使用 confluent cli 的文档。你能建议我怎么做吗!

confluent-platform confluent-cli

5
推荐指数
1
解决办法
372
查看次数

Confluent Cloud - Spring Boot 消费者 REST 端点?

我正在尝试构建一个 Java Spring Boot 应用程序,该应用程序可以发布和获取来自 Confluent Cloud Kafka 的消息。

我按照将Kafka 消息发布到 Confluent Cloud的文章进行操作,并且它有效。

下面是实现

卡夫卡控制器.java

package com.seroter.confluentboot.controller;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

import com.seroter.confluentboot.dto.Product;
import com.seroter.confluentboot.engine.Producer;

@RestController
@RequestMapping(value = "/kafka")
public class KafkaController {

    private final Producer producer;
    
    private final com.seroter.confluentboot.engine.Consumer consumer;

    @Autowired
    KafkaController(Producer producer,com.seroter.confluentboot.engine.Consumer consumer) {
        this.producer = producer;
        this.consumer=consumer;
    }

    @PostMapping(value = "/publish")
    public void sendMessageToKafkaTopic(@RequestParam("message") String message) {
        this.producer.sendMessage(message);
    }
   
    
    @PostMapping(value="/publishJson")
    public ResponseEntity<Product> publishJsonMessage(@RequestBody …
Run Code Online (Sandbox Code Playgroud)

java apache-kafka spring-boot confluent-platform

5
推荐指数
1
解决办法
246
查看次数