如何用Python处理Docker Kafka命令管道输出并预处理后写入日志
问题描述
在终端执行如下Docker命令,会直接将kafka主题的监听结果写入本地的holding_pivot.txt文件:
docker run -it --rm --name consumer --link zookeeper:zookeeper --link kafka:kafka \ debezium/kafka:1.1 watch-topic -a bankserver1.bank.holding > C:/Users/User/python_test/holding_pivot.txt
命令输出的内容结构如下:
WARNING: Using default BROKER_ID=1, which is valid only for non-clustered installations. Using ZOOKEEPER_CONNECT=172.17.0.3:2181 Using KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://172.17.0.6:9092 Using KAFKA_BROKER=172.17.0.4:9092 Contents of topic bankserver1.bank.holding: {"schema":{"type":"struct","fields":[{"type":"struct","fiel... {"schema":{"type":"struct","fields":[{"type":"struct","file...
期望的处理流程为:Docker命令输出 -> Python程序实时处理 -> 写入最终日志文件,避免后续读取文件时区分已处理/未处理行的问题。
调整命令与Python代码后运行无结果,需要实现Python读取管道传入的命令输出的能力。
调整后的错误命令:
# 错误点:换行符使用错误、-it参数不适配管道场景、路径存在多余空格 docker run -it --rm --name consumer --link zookeeper:zookeeper --link / kafka:kafka debezium/kafka:1.1 watch-topic -a bankserver1.bank.holding / | python C:/Users/User/python_test/test_pipe.py > C:/Users/User/python_test / /holding_pivot.txt
编写的错误test_pipe.py代码:
#!/usr/bin/env python3 -u # Note: the -u denotes unbuffered (i.e output straing to stdout without buffering data and then writing to stdout) #!/usr/bin/env python3 import fileinput import json import os import sys from datetime import datetime for line in sys.stdin: with open('log.txt', 'a') as wr: wr.write("Pipe success") with open('log.txt', 'a') as wr: wr.write("Pipe success") with fileinput.input() as f: for line in f: with open('log.txt', 'a') as wr: wr.write(f"Argument List: {str(line)}")
问题原因
- Docker命令的
-it参数中,t参数会分配伪终端,不适配管道传输场景,会导致输出缓冲甚至管道中断 - 命令中的换行符使用错误,Windows环境下cmd的换行符为
^,PowerShell的换行符为`,不能使用/作为换行符 - 命令中的文件路径存在多余空格,导致路径识别错误
- Python代码存在逻辑问题:第一个
for line in sys.stdin是永久阻塞循环,只要流没有中断,后面的代码永远不会执行 - Python代码中使用相对路径
log.txt,写入的路径是执行命令的工作目录,不是Python文件所在目录,容易找不到生成的文件 - 写入文件没有主动刷新缓冲,可能导致内容滞留在内存中没有落地到磁盘
解决方案
可用Python库说明
该需求无需第三方库,Python标准库即可实现:
sys库:直接读取sys.stdin标准输入流,适配管道输入场景fileinput库:也支持读取标准输入流,适合多文件/流混合读取场景
修正后的Docker命令
以Windows PowerShell为例:
docker run -i --rm --name consumer --link zookeeper:zookeeper --link kafka:kafka ` debezium/kafka:1.1 watch-topic -a bankserver1.bank.holding ` | python -u C:/Users/User/python_test/test_pipe.py > C:/Users/User/python_test/holding_pivot.txt
说明:去掉了
t参数,仅保留-i保持标准输入打开;使用-u参数启动Python,强制关闭输出缓冲,保证实时性
修正后的test_pipe.py代码
#!/usr/bin/env python3 import sys import os # 固定log.txt的路径为Python文件所在目录,避免相对路径问题 BASE_DIR = os.path.dirname(os.path.abspath(__file__)) LOG_PATH = os.path.join(BASE_DIR, 'log.txt') def process_line(line): # 这里可以加你自己的处理逻辑,比如过滤JSON行、解析字段等 # 示例:写入调试日志 with open(LOG_PATH, 'a', encoding='utf-8') as f: f.write(f"收到输入行:{line.strip()}\n") # 处理后的结果输出到stdout,最终会写入holding_pivot.txt print(line, end='', flush=True) if __name__ == "__main__": # 逐行读取管道输入 for line in sys.stdin: process_line(line) # 流结束时写入收尾日志 with open(LOG_PATH, 'a', encoding='utf-8') as f: f.write("管道输入流已结束\n")
内容的提问来源于stack exchange,提问作者Chau Loi
相关产品推荐
相关产品推荐

