使用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
排查建议
- 检查Airflow参数注入逻辑:Airflow可能在读取连接参数时自动添加了无效的
timeout参数,librdkafka(Python Kafka客户端依赖)中没有该参数,正确的超时参数为session.timeout.ms、request.timeout.ms等。 - 验证JSON格式正确性:确认Airflow连接
Extra字段中的JSON无语法错误,比如多余逗号,避免解析异常导致参数乱添加。 - 查看完整参数日志:检查Airflow任务执行日志,确认实际传递给Kafka消费者的所有参数,定位
timeout参数来源。 - 适配Hook参数要求:部分Airflow Kafka Hook版本对参数名称有特定要求,确保使用的参数名与Hook期望一致,避免不兼容参数被传递。
内容的提问来源于stack exchange,提问作者Юрий Чернышев
相关产品推荐
相关产品推荐

