如何在Docker Compose中于Kafka Connect就绪后添加自定义连接器?
自动在Kafka Connect就绪后添加自定义连接器的Docker Compose配置方案
核心思路
无需单独维护Shell容器,通过自定义初始化脚本配合Docker Compose的健康检查与依赖控制,实现Kafka Connect服务就绪后自动提交连接器配置。
具体实现步骤
1. 准备连接器配置文件
在项目根目录创建connectors文件夹,将自定义连接器配置保存为JSON文件(例如my-custom-connector.json):
{ "name": "my-custom-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector", "tasks.max": "1", "connection.url": "jdbc:mysql://mysql:3306/mydb", "connection.user": "root", "connection.password": "password", "topic.prefix": "mysql-", "mode": "incrementing", "incrementing.column.name": "id" } }
2. 编写初始化脚本
在项目根目录创建init-connectors.sh脚本,负责等待Kafka Connect就绪并提交配置:
#!/bin/sh # 轮询等待Kafka Connect服务就绪 echo "Waiting for Kafka Connect to be ready..." while ! curl -s http://kafka-connect:8083/connectors > /dev/null; do sleep 5 done echo "Kafka Connect is ready!" # 提交连接器配置 echo "Submitting connector configuration..." curl -X POST http://kafka-connect:8083/connectors \ -H "Content-Type: application/json" \ -d @/connectors/my-custom-connector.json echo "Connector submitted successfully!"
执行命令给脚本添加执行权限:chmod +x init-connectors.sh
3. 修改Docker Compose配置
更新docker-compose.yml,为Kafka Connect添加健康检查,并新增初始化服务:
version: '3.8' services: # 保留原有Zookeeper、Kafka等集群服务配置 zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:latest ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 depends_on: - zookeeper kafka-connect: image: confluentinc/cp-kafka-connect:latest ports: - "8083:8083" environment: CONNECT_BOOTSTRAP_SERVERS: "kafka:9092" CONNECT_REST_ADVERTISED_HOST_NAME: "kafka-connect" CONNECT_GROUP_ID: "connect-group" CONNECT_CONFIG_STORAGE_TOPIC: "connect-configs" CONNECT_OFFSET_STORAGE_TOPIC: "connect-offsets" CONNECT_STATUS_STORAGE_TOPIC: "connect-status" 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_CONFIG_STORAGE_REPLICATION_FACTOR: "1" CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: "1" CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: "1" CONNECT_PLUGIN_PATH: "/usr/share/java,/usr/share/confluent-hub-components" # 添加健康检查,确保服务真正就绪 healthcheck: test: ["CMD", "curl", "-f", "http://localhost:8083/connectors"] interval: 10s timeout: 5s retries: 5 depends_on: - kafka # 新增初始化服务,负责提交连接器配置 init-connectors: image: curlimages/curl:latest volumes: - ./init-connectors.sh:/init-connectors.sh - ./connectors:/connectors command: sh /init-connectors.sh # 依赖Kafka Connect的健康状态,确保就绪后再执行 depends_on: kafka-connect: condition: service_healthy
关键细节说明
- 健康检查:Kafka Connect的
healthcheck通过curl检测服务接口,避免初始化脚本在服务未完全启动时执行。 - 初始化服务:使用轻量的
curlimages/curl镜像,挂载本地脚本和配置文件,无需额外维护Shell容器。 - 依赖控制:通过
depends_on的service_healthy条件,严格保证初始化动作仅在Kafka Connect就绪后触发。
验证方式
启动服务后,可通过以下命令验证:
- 查看初始化日志:
docker-compose logs init-connectors - 检查连接器状态:
curl http://localhost:8083/connectors/my-custom-connector/status
内容的提问来源于stack exchange,提问作者Zikria Azimi
相关产品推荐
相关产品推荐

