Redis Kafka Sink Connector为何写入偏移量字符串?如何禁用?
我正在使用Redis Kafka sink connector从Kafka主题消费消息并更新Redis数据库。消费新消息时,连接器会将消息值以JSON格式upsert到Redis,但同时会额外upsert一个键为com.redis.kafka.connect.sink.$Key的字符串,值格式类似{"topic":"$Key","partition":0,"offset":9}。这种JSON写入符合预期,但额外的字符串写入超出预期。
请问该偏移量字符串写入的原因是什么?是否为连接器记录消费偏移量所必需?如何阻止该行为?
我已查阅Redis Kafka Connector文档、Stack Overflow及Redis论坛帖子,也尝试了GitHub上的Docker示例,但未找到相关解释。
环境配置
Docker Compose文件
version: "2" services: zookeeper: image: quay.io/debezium/zookeeper:2.4 container_name: zookeeper ports: - 2181:2181 - 2888:2888 - 3888:3888 kafka: image: quay.io/debezium/kafka:2.4 container_name: kafka ports: - 29092:29092 links: - zookeeper environment: ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_LISTENERS: INTERNAL://kafka:9092,EXTERNAL://kafka:29092 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:9092,EXTERNAL://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL connect: image: quay.io/debezium/connect:2.4 container_name: connect ports: - 8083:8083 links: - kafka environment: - BOOTSTRAP_SERVERS=kafka:9092 - GROUP_ID=1 - CONFIG_STORAGE_TOPIC=my_connect_configs - OFFSET_STORAGE_TOPIC=my_connect_offsets - STATUS_STORAGE_TOPIC=my_connect_statuses volumes: - ./plugins/redis-redis-kafka-connect-0.9.0:/kafka/connect/redis-connector
Debezium Postgres源连接器配置
{ "name": "postgres-source-connector", "config": { "plugin.name": "pgoutput", "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "tasks.max": "1", "database.hostname": "6.tcp.ngrok.io", "database.port": "18384", "database.user": "postgres", "database.password": "postgres", "database.dbname": "test", "topic.prefix": "dbserver1" } }
Redis sink连接器配置
{ "name": "redis-sink-connector", "config": { "connector.class": "com.redis.kafka.connect.RedisSinkConnector", "tasks.max": "1", "redis.command": "JSONSET", "redis.uri": "redis_uri", "redis.username": "redis_username", "redis.password": "redis_password", "topics": "dbserver1.public.demo", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "true" } }
问题解答
1. 额外写入偏移量键的原因
这个额外的键是Redis Kafka Sink Connector默认启用的偏移量跟踪机制,用于在Redis内部记录每条消息对应的消费偏移量(包含主题、分区、偏移量信息)。它的设计目的是提供一种备选的偏移量存储方式,让连接器可以在某些场景下(比如故障恢复)直接从Redis读取偏移量恢复消费,而不依赖Kafka Connect内置的偏移量存储主题。
2. 是否为消费偏移量记录所必需
不是必需的。Kafka Connect本身已经通过配置中的OFFSET_STORAGE_TOPIC来统一管理所有连接器的消费偏移量,这是官方推荐的标准偏移量存储方案。Redis内部存储的这些偏移量属于连接器额外提供的可选功能,对于使用标准Kafka Connect架构的场景来说,完全可以禁用该功能,不会影响连接器的正常消费与偏移量管理。
3. 如何阻止该行为
在Redis Sink Connector的配置中添加redis.offset.flush.enabled参数并设置为false即可禁用该额外写入行为。修改后的完整配置如下:
{ "name": "redis-sink-connector", "config": { "connector.class": "com.redis.kafka.connect.RedisSinkConnector", "tasks.max": "1", "redis.command": "JSONSET", "redis.uri": "redis_uri", "redis.username": "redis_username", "redis.password": "redis_password", "topics": "dbserver1.public.demo", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "true", "redis.offset.flush.enabled": "false" } }
注:你使用的0.9.0版本连接器已支持该配置项,配置生效后,连接器将不再向Redis写入com.redis.kafka.connect.sink.$Key格式的偏移量键。
内容的提问来源于stack exchange,提问作者Alex Bair

