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

为Python向Kafka Topic推送数据的脚本配置日志求助

为Kafka数据推送脚本添加日志功能

我们直接用Python标准库的logging模块实现日志记录,既能追踪脚本运行状态,也能捕获异常信息。以下是添加日志功能后的完整脚本:

# 导入所需库
import logging
from time import sleep
from json import dumps
from kafka import KafkaProducer
from kafka.errors import KafkaError

# 配置日志:设置级别为INFO,格式包含时间、日志级别、消息内容
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s - %(levelname)s - %(message)s",
    handlers=[
        logging.FileHandler("kafka_producer.log"),  # 日志写入文件
        logging.StreamHandler()  # 同时在控制台输出
    ]
)

def init_kafka_producer():
    """初始化Kafka生产者并返回,失败则记录日志并退出"""
    try:
        producer = KafkaProducer(
            bootstrap_servers=['localhost:9092'],
            value_serializer=lambda x: dumps(x).encode('utf-8')
        )
        logging.info("Kafka生产者初始化成功")
        return producer
    except KafkaError as e:
        logging.error(f"Kafka生产者初始化失败: {str(e)}")
        raise  # 抛出异常终止脚本

if __name__ == "__main__":
    my_producer = init_kafka_producer()

    # 向testnum主题推送数据
    logging.info("开始向testnum主题推送数据")
    for n in range(10):
        my_data = {'num': n}
        try:
            future = my_producer.send('testnum', value=my_data)
            # 等待发送结果,捕获可能的异常
            record_metadata = future.get(timeout=10)
            logging.info(f"成功向testnum主题发送数据: {my_data},分区: {record_metadata.partition},偏移量: {record_metadata.offset}")
        except KafkaError as e:
            logging.error(f"向testnum主题发送数据失败 {my_data}: {str(e)}")
        sleep(1)

    # 向testnum1主题推送偶数数据
    logging.info("开始向testnum1主题推送偶数数据")
    for n in range(10):
        if n % 2 == 0:
            json_data = {'num': n}
            try:
                future = my_producer.send('testnum1', value=json_data)
                record_metadata = future.get(timeout=10)
                logging.info(f"成功向testnum1主题发送数据: {json_data},分区: {record_metadata.partition},偏移量: {record_metadata.offset}")
            except KafkaError as e:
                logging.error(f"向testnum1主题发送数据失败 {json_data}: {str(e)}")
        sleep(1)

    # 关闭生产者并刷新剩余消息
    try:
        my_producer.flush()
        my_producer.close()
        logging.info("Kafka生产者已关闭")
    except KafkaError as e:
        logging.error(f"关闭Kafka生产者失败: {str(e)}")

关键日志功能说明

  • 日志配置:同时将日志输出到控制台和kafka_producer.log文件,兼顾实时查看和事后排查需求
  • 生产者初始化:单独封装成函数,捕获初始化异常并记录错误,避免脚本静默失败
  • 消息发送:通过future.get()等待发送结果,捕获发送过程中的Kafka异常,记录成功/失败的详细信息(包括数据内容、分区、偏移量)
  • 资源清理:关闭生产者前刷新消息,确保所有待发送数据都被处理,并记录关闭状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 09:20:13