标签: confluent-platform

Debezium 的 ExtractNewRecordState 转换无法工作

我正在构建一个数据同步器,它从 MySQL Source 捕获数据更改,并将数据导出到 hive。

我选择使用 Kafka Connect 来实现这一点。我使用 Debezium 作为源连接器,使用 confluence hdfs 作为接收器连接器。

Debezium 提供了单一消息转换,让我可以after从复杂的事件消息中提取字段。我按照列出的文档进行了相同的配置,但它不起作用。

{
    // omit ...
    "transform": "unwrap",
    "transform.unwrap.type": "io.debezium.transforms.ExtractNewRecordState"
}
Run Code Online (Sandbox Code Playgroud)

我尝试在源连接器侧和接收器连接器侧配置转换,但它仍然无法工作。事实上,当我在源连接器端配置它,然后检查相应主题中的消息时,我发现消息仍然包含所有字段,包括beforesource等。

ythh@openstack2:~/confluent-5.5.0$ bin/kafka-avro-console-consumer --from-beginning --bootstrap-server localhost:9092 --topic dbserver1.test_data_1.student3
{"before":null,"after":{"dbserver1.test_data_1.student3.Value":{"id":1,"name":"ggg"}},"source":{"version":"1.1.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1589005572000,"snapshot":{"string":"false"},"db":"test_data_1","table":{"string":"student3"},"server_id":1,"gtid":null,"file":"mysql-bin.000011","pos":9474,"row":0,"thread":{"long":6013},"query":null},"op":"c","ts_ms":{"long":1589005572172},"transaction":null}
{"before":null,"after":{"dbserver1.test_data_1.student3.Value":{"id":2,"name":"no way"}},"source":{"version":"1.1.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1589005893000,"snapshot":{"string":"false"},"db":"test_data_1","table":{"string":"student3"},"server_id":1,"gtid":null,"file":"mysql-bin.000011","pos":11218,"row":0,"thread":{"long":6030},"query":null},"op":"c","ts_ms":{"long":1589005893773},"transaction":null}
{"before":null,"after":{"dbserver1.test_data_1.student3.Value":{"id":3,"name":"not work"}},"source":{"version":"1.1.1.Final","connector":"mysql","name":"dbserver1","ts_ms":1589005900000,"snapshot":{"string":"false"},"db":"test_data_1","table":{"string":"student3"},"server_id":1,"gtid":null,"file":"mysql-bin.000011","pos":11501,"row":0,"thread":{"long":6030},"query":null},"op":"c","ts_ms":{"long":1589005900724},"transaction":null}
Run Code Online (Sandbox Code Playgroud)

我还检查了 kafka 连接日志,这是一些输出:

ythh@openstack2:~/kafka_2.12-2.5.0/logs$ cat connect.log | grep transform
        transforms = []
        transforms = []
        transforms = []
        transforms = []
        transforms = []
        transforms = []
        transforms = []
        transforms = []
        transforms …
Run Code Online (Sandbox Code Playgroud)

apache-kafka apache-kafka-connect debezium confluent-platform

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

无法解决 gradle 中的融合 kafka 依赖关系?

在构建gradle文件中添加汇合kafka的依赖项时,无法解决它。

   compile group: 'io.confluent', name: 'kafka-avro-serializer', version: '4.0.0'
   compile group: 'io.confluent', name: 'kafka-schema-registry', version: '4.0.0'
   compile 'io.confluent:kafka-schema-registry:4.0.0:tests'
Run Code Online (Sandbox Code Playgroud)

添加后出现以下错误。

在此输入图像描述

java apache-kafka confluent-schema-registry confluent-platform

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

卡夫卡满怀期待地关闭了。无法创建许可证主题

在 Kafka 启动时,会记录多条消息kafka/logs/kafkaServer.out并包含:

信息 [代理 0 上的管理员管理器]:处理创建主题请求时出错 CreatableTopic(name='_confluence-license', numPartitions=1,replicationFactor=3,指派=[], configs=[CreateableTopicConfig(name='cleanup.policy', value='compact'), CreateableTopicConfig(name='min.insync.replicas', value='2')]) (kafka.server.AdminManager)

大约 15 分钟后,Kafka 关闭并输出到 kafka/logs/kafkaServer.out

org.apache.kafka.common.errors.InvalidReplicationFactorException: Replication factor: 3 larger than available brokers: 1.
[2020-12-08 04:04:15,951] ERROR [KafkaServer id=0] Fatal error during KafkaServer startup. Prepare to shutdown
Run Code Online (Sandbox Code Playgroud)

(kafka.server.KafkaServer) org.apache.kafka.common.errors.TimeoutException: 无法创建许可证主题 原因: org.apache.kafka.common.errors.InvalidReplicationFactorException: 复制因子: 3 大于可用代理: 1 . [2020-12-08 04:04:15,952] 信息 [KafkaServer id=0] 关闭(kafka.server.KafkaServer)

看来 Kafka 关闭是因为该主题的复制因子设置为 3 _confluent-license?我没有创建主题_confluent-license,这是作为 Kafka 启动许可检查的一部分创建的吗?

为了尝试修复,我进行了修改/v5.5.0/etc/kafka/server.properties,以便内部主题的复制因子为 1:

############################# Internal Topic Settings  #############################
# …
Run Code Online (Sandbox Code Playgroud)

apache-kafka confluent-platform

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

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

如何使用架构注册表/Kafka-Rest 正确注册 Protobuf 架构

我正在尝试使用 kafka-rest 接口将 Protobuf 架构发布到架构注册表:

curl -X POST -H "Content-Type: application/vnd.kafka.protobuf.v2+json" \
   -H "Accept: application/vnd.kafka.v2+json" \
   --data '{"value_schema": "syntax=\"proto3\"; message User { string name = 1; }", "records": [{"value": {"name": "testUser"}}]}' \
   "http://localhost:8082/topics/protobuftest"
Run Code Online (Sandbox Code Playgroud)

我收到此错误:

{"error_code":415,"message":"HTTP 415 Unsupported Media Type"}
Run Code Online (Sandbox Code Playgroud)

问题:指示媒体类型使其发挥作用的正确方法是什么?

protocol-buffers apache-kafka confluent-schema-registry kafka-rest confluent-platform

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

由于 kafka 未定义错误,Docker 构建失败

当我尝试构建我们的 go 应用程序时,我们收到以下错误。

=> ERROR [builder 7/7] RUN CGO_ENABLED=0 GOOS=linux go build -o myapp
 > [builder 7/7] RUN CGO_ENABLED=0 GOOS=linux go build -o myapp:
#14 6.962 # main
#14 6.962 ./kafkaproducer.go:12:12: undefined: kafka.NewProducer
#14 6.962 ./kafkaproducer.go:12:31: undefined: kafka.ConfigMap
#14 6.962 ./kafkaproducer.go:23:10: undefined: kafka.Message
#14 6.962 ./kafkaproducer.go:39:13: undefined: kafka.Message
Run Code Online (Sandbox Code Playgroud)

我的 Docker 文件是

FROM golang:1.16-alpine AS builder
RUN mkdir /app
ADD . /app
WORKDIR /app
RUN go mod tidy
RUN CGO_ENABLED=0 GOOS=linux go build -o myapp
FROM busybox AS prod
COPY --from=builder …
Run Code Online (Sandbox Code Playgroud)

go apache-kafka docker dockerfile confluent-platform

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

在Kafka中执行批量验证并发送到相应的主题

我在 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是三个不同的数据。

这是扫描仪:

  1. 在将批次结束记录推送到B1E主题之前,我们经过一些验证后得知该B1批次无效。

  2. 在这种情况下,B1批次范围内BS的所有数据都BE应该转到特定主题。

因此,在上面的示例中,b1批次应转到主题 T1,b2批次应转到主题 T2。

我怎样才能使用卡夫卡做到这一点?

apache-kafka confluent-platform

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

LibrdKafkaError:经纪人:随机运行约 2 小时后未知成员

现在,我想实施node-rdkafka到我们的服务中,但我多次遇到此错误Broker: Unknown member。github 上的同一问题是https://github.com/confluenceinc/confluence-kafka-dotnet/issues/1464。他们说我们的消费者使用相同的组 ID 来重试或延迟。但我没有发现我的代码有任何重试和延迟。或https://github.com/confluenceinc/confluence-kafka-python/issues/1004,但我重新检查了所有消费者组 ID,它是唯一的。

生产者的配置node-rdkafka如下:

      this.producer = new Producer({
        "client.id": this.cliendID,
        "metadata.broker.list": this.brokerList,
        'compression.codec': "lz4",
        'retry.backoff.ms': 200,
        'socket.keepalive.enable': true,
        'queue.buffering.max.messages': 100000,
        'queue.buffering.max.ms': 1000,
        'batch.num.messages': 1000000,
        "transaction.timeout.ms": 2000,
        "enable.idempotence": false,
        "max.in.flight.requests.per.connection": 1,
        "debug": this.debug,
        'dr_cb': true,
        "retries": 0,
        "log_cb": (_: any) => console.log(`log_cb =>`, _),
        "sasl.username": this.saslUsername,
        "sasl.password": this.saslPassword,
        "sasl.mechanism": this.saslMechanism,
        "security.protocol": this.securityProtocol
      }, {
        "acks": -1
      })
Run Code Online (Sandbox Code Playgroud)

Consumer的配置node-rdkafka如下:

this.consumer = new KafkaConsumer({
        'group.id': this.groupID,
        'metadata.broker.list': …
Run Code Online (Sandbox Code Playgroud)

apache-kafka librdkafka kafkajs confluent-platform

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

带有自定义时间戳的 Kafka Connect.extractor

我在尝试将消息从 Kafka 读取到 S3 时遇到了将 jar 添加到 Kafka 连接类路径的问题。

目标是根据时间戳在分区中写入消息,时间戳是 Kafka 消息中 Key 的一部分。

为了使故事简短,我必须提供自定义时间戳提取器。按照此处的文档创建了一个实现TimestampExtractor接口的类并将 JAR 位置添加到plugin.path属性中。

问题是当我开始连接时,找不到类。不知何故,jar 不在类路径中,我得到了

org.apache.kafka.common.config.ConfigException: Invalid timestamp extractor: partitioner.SpotadDateTimeExtractor
Run Code Online (Sandbox Code Playgroud)

附加数据:

版本:融合 4.0.0

连接:连接独立

启动命令:

sudo /home/ubuntu/confluent-4.0.0/bin/connect-standalone \ /home/ubuntu/confluent-4.0.0/etc/kafka/connect-standalone.properties \ /home/ubuntu/confluent-4.0.0/etc/kafka-connect-s3/quickstart-s3.properties

Apreaciate任何帮助。

apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
991
查看次数

如何正确重启 kafka s3 接收器连接?

自 5 月 1 日起,我开始使用 kafka s3 接收器连接器(来自融合包的捆绑连接器)。它可以正常工作到 5 月 8 日。检查状态,它告诉某些 aws 异常使此连接器崩溃。这应该不是什么大问题,所以我想恢复它。

我尝试了以下步骤:

  1. 我 POST /connectors/s3sink/restart 。然后我看到连接器处于 RUNNING 模式,但任务仍然 FAIL。
  2. 然后我 PUT /connectors/s3sink/task/0/restart。好的,现在任务处于 RUNNING 模式。

但是后来我跟踪日志,发现它开始重写旧数据,例如5月3日的数据。它弄乱了旧数据!

那么,connect restart REST API 是否会重置偏移量?我认为它会保存偏移量并从它失败的偏移量开始。

以及如何正确重启失败的连接器任务?通过删除那些 POD?(使用 kubernetes),还是通过 REST /task/0/restart?我什么时候应该使用/connectors/s3sink/restart?

amazon-s3 apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
4333
查看次数