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

Python使用kafka-logging-handler通过YAML配置推送日志到Kafka报错求助

问题原因

你遇到的报错是配置文件的低级拼写错误+参数配置不当导致的,具体问题如下:

  1. handler命名与引用不匹配
    你在handlers配置块中定义的Kafka handler名称为kakfa_log_hander,存在两处拼写错误:
  • 单词kafka拼写错误,将kafka写成了kakfa(f和k顺序颠倒)
  • 单词handler拼写错误,将handler写成了hander(缺少字母l)
    而你在kfa_logger的handlers列表中引用的是kafka_log_handler,名称完全不匹配,导致dictConfig加载时找不到对应的handler定义,抛出无法添加handler的错误。
  1. Kafka连接参数硬编码
    当前配置中hosts_list写死为MY_IP_ADDRESS:9091,没有从环境变量读取,即便修复拼写错误也会出现Kafka连接失败的问题。

  2. 默认日志无Kafka handler绑定
    你测试时使用logging.info("test")输出日志,默认调用的是root logger,但你的配置文件中仅定义了三个自定义logger,没有配置root logger的handler,所以默认日志不会推送到Kafka。

修复方案

替换你的logconfig.yml为以下内容即可:

version: 1
formatters:
  default:
    format: '%(threadName)s:[%(asctime)s] %(levelname)s in %(module)s: %(message)s'
  json:
    class: pythonjsonlogger.jsonlogger.JsonFormatter
    format: '%(threadName)s:[%(asctime)s] %(levelname)s in %(module)s: %(message)s'
handlers:
  console_handler:
    level: INFO
    class: logging.StreamHandler
    formatter: json
    stream: ext://sys.stdout
  # 已修正handler名称拼写
  kafka_log_handler:
    level: INFO
    class: kafka_logger.handlers.KafkaLoggingHandler
    # 从环境变量读取Kafka连接地址
    hosts_list: ext://os.environ.BOOTSTRAP_SERVERS
    # 从环境变量读取Topic名称
    topic: ext://os.environ.TOPIC_NAME
    security_protocol: "PLAINTEXT"
    kafka_producer_args:
      api_version_auto_timeout_ms: 1000000
      request_timeout_ms: 1000000
    additional_fields:
      service: test-service
    formatter: json
loggers:
  default_logger:
    level: INFO
    handlers: [console_handler]
  debug_logger:
    level: DEBUG
    handlers: [console_handler]
  kfa_logger:
    level: INFO
    handlers: [console_handler, kafka_log_handler]
# 新增root logger配置,让默认日志输出也能推送到Kafka
root:
  level: INFO
  handlers: [console_handler, kafka_log_handler]
disable_existing_loggers: False

如果你不想用root logger输出,也可以在测试代码中指定使用自定义的kfa_logger:

if __name__ == "__main__":
    # 替换原有的logging.info调用
    logger = logging.getLogger("kfa_logger")
    logger.info("test")
    
    # 剩余原有代码不变
    msg = str.encode("test message")
    producer = KafkaProducer(bootstrap_servers=os.environ.get("BOOTSTRAP_SERVERS"))
    producer.send(os.environ.get("TOPIC_NAME"), msg)
    producer.flush()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 12:27:04