【发布时间】:2021-05-26 15:04:02
【问题描述】:
说明
- 我在 Kafka Connect 分布式模式下有一个 pubSubSource 连接器,它只是从 PubSub 订阅中读取数据并写入 Kafka 主题。问题是,即使我将一条消息发布到 GCP PubSub,我也会在我的 Kafka 主题中两次写入这条消息。
如何重现
-
部署 Kafka 和 Kafka 连接
-
使用以下
pubSubSource配置创建连接器:curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{ "name": "pubSubSource", "config": { "connector.class":"com.google.pubsub.kafka.source.CloudPubSubSourceConnector", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter", "tasks.max":"1", "cps.subscription":"pubsub-test-sub", "kafka.topic":"kafka-sub-topic", "cps.project":"test-project123", "gcp.credentials.file.path":"/tmp/gcp-creds/account-key.json" } }' -
以下是 Kafka 连接配置:
"plugin.path": "/usr/share/java,/usr/share/confluent-hub-components" "key.converter": "org.apache.kafka.connect.json.JsonConverter" "value.converter": "org.apache.kafka.connect.json.JsonConverter" "key.converter.schemas.enable": "false" "value.converter.schemas.enable": "false" "internal.key.converter": "org.apache.kafka.connect.json.JsonConverter" "internal.value.converter": "org.apache.kafka.connect.json.JsonConverter" "config.storage.replication.factor": "1" "offset.storage.replication.factor": "1" "status.storage.replication.factor": "1" -
使用以下命令向 PubSub 主题发布消息:
gcloud pubsub topics publish test-topic --message='{"someKey":"someValue"}' -
从目标 Kafka 主题读取消息:
/usr/bin/kafka-console-consumer --bootstrap-server xx.xxx.xxx.xx:9092 --topic kafka-topic --from-beginning # Output {"someKey":"someValue"} {"someKey":"someValue"}
为什么会这样,是不是我做错了什么?
【问题讨论】:
标签: apache-kafka apache-kafka-connect google-cloud-pubsub confluent-platform google-cloud-pubsublite