我正在构建一个数据同步器,它从 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)
我尝试在源连接器侧和接收器连接器侧配置转换,但它仍然无法工作。事实上,当我在源连接器端配置它,然后检查相应主题中的消息时,我发现消息仍然包含所有字段,包括before、source等。
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
在构建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
在 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:
Run Code Online (Sandbox Code Playgroud)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(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) 查看https://github.com/confluenceinc/confluence-kafka-dotnet/blob/master/examples/JsonSerialization/Program.cs它需要架构注册表 URL。有没有一种简单的方法来序列化/反序列化 JSON,而不需要额外的复杂性?
我正在尝试使用 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
当我尝试构建我们的 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) 我在 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。
我怎样才能使用卡夫卡做到这一点?
现在,我想实施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) 我在尝试将消息从 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任何帮助。
自 5 月 1 日起,我开始使用 kafka s3 接收器连接器(来自融合包的捆绑连接器)。它可以正常工作到 5 月 8 日。检查状态,它告诉某些 aws 异常使此连接器崩溃。这应该不是什么大问题,所以我想恢复它。
我尝试了以下步骤:
但是后来我跟踪日志,发现它开始重写旧数据,例如5月3日的数据。它弄乱了旧数据!
那么,connect restart REST API 是否会重置偏移量?我认为它会保存偏移量并从它失败的偏移量开始。
以及如何正确重启失败的连接器任务?通过删除那些 POD?(使用 kubernetes),还是通过 REST /task/0/restart?我什么时候应该使用/connectors/s3sink/restart?
amazon-s3 apache-kafka apache-kafka-connect confluent-platform
apache-kafka ×10
.net ×1
.net-core ×1
amazon-s3 ×1
debezium ×1
docker ×1
dockerfile ×1
go ×1
java ×1
json ×1
kafka-rest ×1
kafkajs ×1
librdkafka ×1