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

使用Airflow连接Kafka时出现配置错误求助

Airflow Kafka连接配置报错排查请求

问题现象

本地搭建Airflow与Kafka环境,Python显式传入连接参数创建消费者时一切正常,但将参数配置到Airflow连接中时,触发以下错误:

KafkaError{code=_INVALID_ARG,val=-186,str="No such configuration property: "timeout""}

使用的连接参数

{  
  "bootstrap.servers": "kafka:9092",  
  "group.id": "kafka_11",  
  "auto.offset.reset": "earliest"
}

环境说明

  • Airflow与Kafka运行在不同Docker容器,已加入同一网络
  • 单独连接Kafka、配置Kafka Connect均正常

相关配置文件

Airflow Docker配置

使用官方2.10.2版本的docker-compose.yaml

Kafka Docker配置

services:
  akhq:
    image: tchiotludo/akhq
    container_name: AKHQ
    restart: always
    environment:
      AKHQ_CONFIGURATION: |
        akhq:
          connections:
            docker-kafka-server:
              properties:
                bootstrap.servers: "kafka:9092"
              schema-registry:
                url: "http://schema-registry:8085"
              connect:
                - name: "connect"
                  url: "http://connect:8083"
    ports:
      - 8081:8080
    links:
      - kafka
      - schema-registry
    networks:
      - my_network

  zookeeper:
    image: confluentinc/cp-zookeeper:${CONFLUENT_VERSION:-latest}
    container_name: zookeeper
    restart: always
    volumes:
      - zookeeper-data:/var/lib/zookeeper/data:Z
      - zookeeper-log:/var/lib/zookeeper/log:Z
    environment:
      ZOOKEEPER_CLIENT_PORT: '2181'
      ZOOKEEPER_ADMIN_ENABLE_SERVER: 'false'
    networks:
      - my_network

  kafka:
    image: confluentinc/cp-kafka:${CONFLUENT_VERSION:-latest}
    container_name: kafka
    restart: always
    volumes:
      - kafka-data:/var/lib/kafka/data:Z
    environment:
        KAFKA_BROKER_ID: '0'
        KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
        KAFKA_NUM_PARTITIONS: '12'
        KAFKA_COMPRESSION_TYPE: 'gzip'
        KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: '1'
        KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: '1'
        KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: '1'
        KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://kafka:9092'
        KAFKA_CONFLUENT_SUPPORT_METRICS_ENABLE: 'false'
        KAFKA_JMX_PORT: '9091'
        KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
        KAFKA_AUTHORIZER_CLASS_NAME: 'kafka.security.authorizer.AclAuthorizer'
        KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: 'true'
    links:
      - zookeeper
    ports:
      - 9092:9092
    networks:
      - my_network

  schema-registry:
    image: confluentinc/cp-schema-registry:${CONFLUENT_VERSION:-latest}
    container_name: schema-registry
    restart: always
    depends_on:
      - kafka
    environment:
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'PLAINTEXT://kafka:9092'
      SCHEMA_REGISTRY_HOST_NAME: 'schema-registry'
      SCHEMA_REGISTRY_LISTENERS: 'http://0.0.0.0:8085'
      SCHEMA_REGISTRY_LOG4J_ROOT_LOGLEVEL: 'INFO'
    networks:
      - my_network

  connect:
    image: confluentinc/cp-kafka-connect:${CONFLUENT_VERSION:-latest}
    container_name: kafka-connect
    restart: always
    depends_on:
      - kafka
      - schema-registry
    environment:
      CONNECT_BOOTSTRAP_SERVERS: 'kafka:9092'
      CONNECT_REST_PORT: '8083'
      CONNECT_REST_LISTENERS: 'http://0.0.0.0:8083'
      CONNECT_REST_ADVERTISED_HOST_NAME: 'connect'
      CONNECT_CONFIG_STORAGE_TOPIC: '__connect-config'
      CONNECT_OFFSET_STORAGE_TOPIC: '__connect-offsets'
      CONNECT_STATUS_STORAGE_TOPIC: '__connect-status'
      CONNECT_GROUP_ID: 'kafka-connect'
      CONNECT_KEY_CONVERTER_SCHEMAS_ENABLE: 'true'
      CONNECT_KEY_CONVERTER: 'io.confluent.connect.avro.AvroConverter'
      CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL: 'http://schema-registry:8085'
      CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: 'true'
      CONNECT_VALUE_CONVERTER: 'io.confluent.connect.avro.AvroConverter'
      CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: 'http://schema-registry:8085'
      CONNECT_INTERNAL_KEY_CONVERTER: 'org.apache.kafka.connect.json.JsonConverter'
      CONNECT_INTERNAL_VALUE_CONVERTER: 'org.apache.kafka.connect.json.JsonConverter'
      CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: '1'
      CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: '1'
      CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: '1'
      CONNECT_PLUGIN_PATH: '/usr/share/java/,/usr/share/confluent-hub-components/'
    volumes:
      - ./data/kafka_properties/connect-distributed.properties:/etc/kafka/connect-distributed.properties
      - ./data/kafka_properties/connect-standalone.properties:/etc/kafka/connect-standalone.properties
      - ./data/kafka_components:/usr/share/confluent-hub-components
    ports:
      - 8083:8083
      - 9292:9092
    networks:
      - my_network

  minio:
    image: minio/minio:RELEASE.2024-07-04T14-25-45Z
    restart: always
    volumes:
      - ./data:/data
    environment:
      - MINIO_ROOT_USER=admin
      - MINIO_ROOT_PASSWORD=admin
    command: server /data --console-address ":9002"
    ports:
      - "9000:9000"
      - "9002:9002"
    networks:
      - my_network

networks:
  my_network:
    external: true

排查建议

  1. 检查Airflow参数注入逻辑:Airflow可能在读取连接参数时自动添加了无效的timeout参数,librdkafka(Python Kafka客户端依赖)中没有该参数,正确的超时参数为session.timeout.ms、request.timeout.ms等。
  2. 验证JSON格式正确性:确认Airflow连接Extra字段中的JSON无语法错误,比如多余逗号,避免解析异常导致参数乱添加。
  3. 查看完整参数日志:检查Airflow任务执行日志,确认实际传递给Kafka消费者的所有参数,定位timeout参数来源。
  4. 适配Hook参数要求:部分Airflow Kafka Hook版本对参数名称有特定要求,确保使用的参数名与Hook期望一致,避免不兼容参数被传递。

内容的提问来源于stack exchange,提问作者Юрий Чернышев

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 11:24:50