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

Docker-Compose部署Filebeat-Kafka-Clickhouse:Kafka无法写入Clickhouse排查

问题排查:Filebeat-Kafka-Clickhouse链路中Kafka到Clickhouse数据写入失败

问题描述

通过Docker-Compose搭建Filebeat-Kafka-Clickhouse日志链路,Filebeat收集日志发送至Kafka正常,但Kafka无法将日志写入Clickhouse,容器日志无异常。

配置文件

docker-compose.yml

version: "3.5"
services:
  filebeat:
    image: docker.elastic.co/beats/filebeat:8.0.1
    volumes:
      - "./logs:/var/log/"
      - "./filebeat.yml:/usr/share/filebeat/filebeat.yml"
    container_name: filebeat
    networks:
      - filebeat-net
    depends_on:
      - kafka
  zookeeper:
    image: docker.io/bitnami/zookeeper:3.7
    ports:
      - "2181:2181"
    volumes:
      - "zookeeper-data:/bitnami"
    environment:
      - ALLOW_ANONYMOUS_LOGIN=yes
    networks:
      - filebeat-net
    container_name: zookeeper
  kafka:
    image: docker.io/bitnami/kafka:3
    ports:
      - "9092:9092"
      - '29092:29092'
    volumes:
      - "kafka-data:/bitnami"
    environment:
      - "HOSTNAME_COMMAND=docker info | grep ^Name: | cut -d' ' -f 2"
      - "KAFKA_CFG_INTER_BROKER_LISTENER_NAME=PLAINTEXT"
      - "ALLOW_PLAINTEXT_LISTENER=yes"
      - "KAFKA_CFG_ZOOKEEPER_CONNECT=zookeeper:2181"
      - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092
      - "KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT"
      - "KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,PLAINTEXT_HOST://:29092"
      - "KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092"
      - "KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true"
    depends_on:
      - zookeeper
    networks:
      - filebeat-net  
    container_name: kafka

  clickhouse-server:
    image: clickhouse/clickhouse-server:22.5.1
    container_name: clickhouse-server
    hostname: clickhouse-server
    ulimits:
      nofile:
        soft: 262144
        hard: 262144
    ports:
      - "8123:8123"
      - "9000:9000"
    volumes:
      - './init.sql:/docker-entrypoint-initdb.d/init-db.sql'
    depends_on:
      - kafka
    networks:
      - filebeat-net


networks:
  filebeat-net:
    name: filebeat-net
    driver: bridge

volumes:
  zookeeper-data:
    name: zookeeper-data
    driver: local
  kafka-data:
    name: kafka-data
    driver: local

filebeat.yml

filebeat.inputs:
- input_type: log
  paths:
  - /var/log/*.log

output.kafka:
  hosts: ["kafka:9092"]
  topic: hello-messages

init.sql

CREATE TABLE IF NOT EXISTS mylogger (  
     id String, 
          event_time DateTime64(6),  
            details_json String)  
            ENGINE  = MergeTree() PARTITION BY toYYYYMM(event_time) ORDER BY (id, event_time) SETTINGS index_granularity=8192;

CREATE TABLE  IF NOT EXISTS 
mylogger_kafka ( payload String ) ENGINE = Kafka('kafka:9092', 'hello-messages', 'KAFKA2CH_click', 'JSONAsString');   


CREATE MATERIALIZED VIEW IF NOT EXISTS mylogger_kafka_consumer TO mylogger AS 
SELECT JSONExtractString(payload, 'payload', 'after', 'id')    as id, 
 toDateTime64(JSONExtractString(payload, 'payload', 'after', 'event_time'), 3, 'Asia/Jerusalem') as event_time, 
 JSONExtractString(payload, 'payload', 'after', 'details_json')    as details_json 
 FROM mylogger_kafka;

排查要点及修复方案

1. 核心问题:Filebeat输出结构与Clickhouse解析逻辑不匹配

Filebeat默认发送到Kafka的日志结构并不包含payload.after嵌套字段,实际格式类似:

{"@timestamp":"2024-05-20T10:00:00.000Z","message":"原始日志内容","host":{"name":"filebeat"},"log":{"file":{"path":"/var/log/test.log"}}}

而你的物化视图中使用JSONExtractString(payload, 'payload', 'after', 'id')解析,完全不匹配实际消息结构,导致无法提取数据。

修复步骤:

  • 修改init.sql中的Kafka表定义,将格式从JSONAsString改为JSON,让Clickhouse自动解析JSON字段:
    CREATE TABLE IF NOT EXISTS mylogger_kafka (
        `@timestamp` DateTime64,
        message String,
        host Map(String, String)
    ) ENGINE = Kafka('kafka:9092', 'hello-messages', 'KAFKA2CH_click', 'JSON');
    
  • 调整物化视图的解析逻辑,适配Filebeat的输出结构:
    CREATE MATERIALIZED VIEW IF NOT EXISTS mylogger_kafka_consumer TO mylogger AS 
    SELECT
        generateUUIDv4() AS id, -- 生成唯一ID,或根据日志内容提取
        toDateTime64(`@timestamp`, 3, 'Asia/Jerusalem') AS event_time,
        message AS details_json
    FROM mylogger_kafka;
    

2. 检查Kafka消息内容确认结构

进入Kafka容器,查看实际消息内容,验证结构是否符合预期:

docker exec -it kafka /opt/bitnami/kafka/bin/kafka-console-consumer.sh --bootstrap-server kafka:9092 --topic hello-messages --from-beginning

3. 重置Kafka消费者偏移量

如果消费者组KAFKA2CH_click已经存在旧的偏移量,可能导致新消息不被消费,执行以下命令重置偏移量:

# 进入Clickhouse客户端
docker exec -it clickhouse-server clickhouse-client
# 重置偏移量到最早位置
ALTER TABLE mylogger_kafka RESET KAFKA OFFSET TO EARLIEST;

4. 确认初始化脚本执行状态

查看Clickhouse容器日志,确认init-db.sql是否被正确执行:

docker logs clickhouse-server | grep "init-db.sql"

如果脚本未执行,可手动进入Clickhouse客户端执行init.sql中的语句。

5. 验证网络连通性

确认Clickhouse容器能正常访问Kafka服务:

docker exec -it clickhouse-server ping kafka

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 08:56:08