You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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就绪后触发。

验证方式

启动服务后,可通过以下命令验证:

  1. 查看初始化日志:docker-compose logs init-connectors
  2. 检查连接器状态:curl http://localhost:8083/connectors/my-custom-connector/status

内容的提问来源于stack exchange,提问作者Zikria Azimi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 05:15:33