Kafka 集群部署与 Connect 同步到 Elasticsearch
August 20, 2024About 2 min
Kafka 集群部署与 Connect 同步到 Elasticsearch
用 docker-compose 起一套 ZooKeeper 架构的三节点 Kafka 集群,再配 Kafka Connect 把 topic 数据自动同步到 Elasticsearch,一个 topic 对应一个 index。下面把 IP 都写成 <HOST_IP> 占位,部署时换成宿主机地址。
三节点集群(ZooKeeper 架构)
用 bitnami 的镜像,一个 ZooKeeper 控制器 + 三个 broker + 一个可视化管理页。关键在每个 broker 的 KAFKA_CFG_ADVERTISED_LISTENERS 要写宿主机能访问到的地址,否则外部客户端连不上。
version: "2"
networks:
kafka-network:
driver: bridge
services:
zookeeper:
image: docker.io/bitnami/zookeeper:3.9
container_name: zookeeper
networks:
- kafka-network
ports:
- "2181:2181"
volumes:
- "zookeeper_data:/bitnami"
environment:
- ALLOW_ANONYMOUS_LOGIN=yes
restart: unless-stopped
kafka-0:
image: docker.io/bitnami/kafka:3.6
container_name: kafka-0
networks:
- kafka-network
ports:
- '9092:9092'
environment:
- KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_BROKER_ID=0
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9092
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://<HOST_IP>:9092
- ALLOW_PLAINTEXT_LISTENER=yes
volumes:
- kafka_0_data:/bitnami/kafka
depends_on:
- zookeeper
restart: unless-stopped
kafka-1:
image: docker.io/bitnami/kafka:3.6
container_name: kafka-1
networks:
- kafka-network
ports:
- '9093:9093'
environment:
- KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_BROKER_ID=1
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9093
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://<HOST_IP>:9093
- ALLOW_PLAINTEXT_LISTENER=yes
volumes:
- kafka_1_data:/bitnami/kafka
depends_on:
- zookeeper
restart: unless-stopped
kafka-2:
image: docker.io/bitnami/kafka:3.6
container_name: kafka-2
networks:
- kafka-network
ports:
- '9094:9094'
environment:
- KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_BROKER_ID=2
- KAFKA_CFG_LISTENERS=PLAINTEXT://:9094
- KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://<HOST_IP>:9094
- ALLOW_PLAINTEXT_LISTENER=yes
volumes:
- kafka_2_data:/bitnami/kafka
depends_on:
- zookeeper
restart: unless-stopped
kafka_manager:
image: hlebalbau/kafka-manager:stable
container_name: kafka_manager
networks:
- kafka-network
ports:
- "9000:9000"
environment:
ZK_HOSTS: "zookeeper:2181"
APPLICATION_SECRET: "random-secret"
depends_on:
- zookeeper
- kafka-0
- kafka-1
- kafka-2
restart: unless-stopped
volumes:
kafka_0_data:
driver: local
kafka_1_data:
driver: local
kafka_2_data:
driver: local
zookeeper_data:
driver: local
生产者 / 消费者测试
进 kafka 容器,建 topic 并收发消息:
cd /opt/bitnami/kafka/config
# 生产者
kafka-console-producer.sh --bootstrap-server <HOST_IP>:9092,<HOST_IP>:9093,<HOST_IP>:9094 --topic test
# 消费者,--from-beginning 从头读
kafka-console-consumer.sh --bootstrap-server <HOST_IP>:9092,<HOST_IP>:9093,<HOST_IP>:9094 --topic test --from-beginning
Kafka Connect 同步到 Elasticsearch
目标:自动把 Kafka topic 的数据发到 Elasticsearch,以 topic 开头的 topic 都同步,一个 topic 对应一个 index。
基于 cp-kafka-connect 镜像装 ES connector 插件,并用 sed 改掉默认地址。把 ES 用户名密码作为占位,build 时按需替换:
FROM confluentinc/cp-kafka-connect:7.6.0
# 装 elasticsearch connector 插件
RUN confluent-hub install --no-prompt confluentinc/kafka-connect-elasticsearch:14.0.12
# 改 kafka 地址
RUN sed -i 's/localhost:9092/<HOST_IP>:9092,<HOST_IP>:9093,<HOST_IP>:9094/g' /etc/kafka/connect-standalone.properties && \
# 禁用 schemas
sed -i 's/key.converter.schemas.enable=true/key.converter.schemas.enable=false/g' /etc/kafka/connect-standalone.properties && \
sed -i 's/value.converter.schemas.enable=true/value.converter.schemas.enable=false/g' /etc/kafka/connect-standalone.properties && \
# 改 elasticsearch 地址
sed -i 's/localhost:9200/<HOST_IP>:9200/g' /usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties && \
# 匹配以 topic 开头的 topic
sed -i 's/topics=test-elasticsearch-sink/topics.regex=topic.*/g' /usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties && \
echo schema.ignore=true >> /usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties && \
# ES 用户名密码
echo connection.username=elastic >> /usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties && \
echo connection.password=<YOUR_PASSWORD> >> /usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties
build:
docker build -t elastic-connector:test .
启动 connector(--net=host 用宿主机网络):
sudo docker run -d \
--name=kafka-connect \
--net=host \
-e CONNECT_BOOTSTRAP_SERVERS=<HOST_IP>:9092,<HOST_IP>:9093,<HOST_IP>:9094 \
-e CONNECT_REST_PORT=8083 \
-e CONNECT_GROUP_ID="kafka-connect" \
-e CONNECT_CONFIG_STORAGE_TOPIC="kafka-connect-config" \
-e CONNECT_OFFSET_STORAGE_TOPIC="kafka-connect-offsets" \
-e CONNECT_STATUS_STORAGE_TOPIC="kafka-connect-status" \
-e CONNECT_KEY_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \
-e CONNECT_VALUE_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \
-e CONNECT_INTERNAL_KEY_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \
-e CONNECT_INTERNAL_VALUE_CONVERTER="org.apache.kafka.connect.json.JsonConverter" \
-e CONNECT_REST_ADVERTISED_HOST_NAME="localhost" \
-e CONNECT_PLUGIN_PATH=/usr/share/java,/usr/share/confluent-hub-components \
elastic-connector:test /bin/connect-standalone /etc/kafka/connect-standalone.properties \
/usr/share/confluent-hub-components/confluentinc-kafka-connect-elasticsearch/etc/quickstart-elasticsearch.properties