Docker Compose部署纯Kafka及Kafka Connect:ZooKeeper兼容问题排查
Kafka ZooKeeper模式启动修复与PostgreSQL Sink Connector配置
一、解决Kafka ZooKeeper与KRaft启动冲突问题
你的问题根源是Kafka配置中同时混合了ZK模式和KRaft模式的参数,导致启动逻辑冲突:
- 若配置了
process.roles、node.id等KRaft专属参数,Kafka会强制要求使用KRaft模式,忽略ZK配置; - 若删除KRaft参数但未正确配置
zookeeper.connect,又会触发ZK依赖异常。
调整后的Docker Compose配置(纯ZK模式)
删除所有KRaft相关配置项,确保Kafka仅依赖ZK启动:
version: '3.8' services: zookeeper: image: apache/kafka:latest command: ["zookeeper-server-start.sh", "/opt/kafka/config/zookeeper.properties"] ports: - "2181:2181" volumes: - zk-data:/tmp/zookeeper kafka: image: apache/kafka:latest command: ["kafka-server-start.sh", "/opt/kafka/config/server.properties"] ports: - "9092:9092" environment: # 核心ZK连接配置 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 # 监听与通告地址配置 KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:9092 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 # 单节点环境必备配置 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 volumes: - kafka-data:/tmp/kafka-logs depends_on: - zookeeper volumes: zk-data: kafka-data:
验证启动状态
启动后执行以下命令确认Kafka正常运行:
# 列出当前所有topic docker exec -it <kafka-container-id> kafka-topics.sh --list --bootstrap-server localhost:9092
二、PostgreSQL Sink Connector配置(非Confluent方案)
使用Apache Kafka Connect的JDBC Sink Connector(社区开源,无专有依赖)实现Kafka消息到PostgreSQL的同步。
1. 启用Kafka Connect服务
在上述Docker Compose中添加Kafka Connect服务:
kafka-connect: image: apache/kafka:latest command: ["connect-distributed.sh", "/opt/kafka/config/connect-distributed.properties"] ports: - "8083:8083" environment: CONNECT_BOOTSTRAP_SERVERS: kafka:9092 CONNECT_GROUP_ID: connect-group # 存储Connector配置的topic(需确保Kafka已创建这些topic,或让Connect自动创建) CONNECT_CONFIG_STORAGE_TOPIC: connect-configs CONNECT_OFFSET_STORAGE_TOPIC: connect-offsets CONNECT_STATUS_STORAGE_TOPIC: connect-status # JSON格式转换器(可根据消息格式调整为Avro等) CONNECT_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter CONNECT_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter CONNECT_INTERNAL_KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter CONNECT_INTERNAL_VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter CONNECT_REST_ADVERTISED_HOST_NAME: kafka-connect # 插件目录,用于存放JDBC驱动 CONNECT_PLUGIN_PATH: "/opt/kafka/plugins" volumes: # 挂载PostgreSQL JDBC驱动到插件目录(需提前下载对应版本的jar包) - ./postgresql-42.7.3.jar:/opt/kafka/plugins/postgresql-42.7.3.jar depends_on: - kafka
2. 下载PostgreSQL JDBC驱动
从PostgreSQL官方下载对应版本的JDBC驱动jar包(如postgresql-42.7.3.jar),放到Docker Compose所在目录,确保挂载路径正确。
3. 创建Sink Connector
通过Connect的REST API提交配置:
curl -X POST -H "Content-Type: application/json" http://localhost:8083/connectors -d '{ "name": "postgres-sink-connector", "config": { "connector.class": "org.apache.kafka.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", # 要同步的Kafka topic,多个用逗号分隔 "topics": "your-target-topic", # PostgreSQL连接信息 "connection.url": "jdbc:postgresql://postgres-host:5432/your-database", "connection.user": "your-db-user", "connection.password": "your-db-password", # 自动创建表(首次同步时启用) "auto.create": "true", # 自动更新表结构(当消息字段变化时启用) "auto.evolve": "true", # 插入模式:upsert表示存在则更新,不存在则插入 "insert.mode": "upsert", # 主键字段(需与消息中的字段对应) "pk.fields": "id", # 主键取值来源:record_value表示从消息内容中获取 "pk.mode": "record_value" } }'
4. 验证同步效果
- 向目标Kafka topic发送一条测试消息:
docker exec -it <kafka-container-id> kafka-console-producer.sh --broker-list localhost:9092 --topic your-target-topic # 输入测试消息,例如:{"id":1,"content":"test message"}
- 登录PostgreSQL数据库,查看对应表是否已生成数据。
内容的提问来源于stack exchange,提问作者Gleichmut
相关产品推荐
相关产品推荐

