【发布时间】:2022-01-03 05:26:11
【问题描述】:
请花时间阅读完整的问题以了解确切的问题。谢谢。
我有一个运行程序/驱动程序,它监听 Kafka 主题并在收到有关该主题的新消息时使用 ThreadPoolExecuter 调度任务(如下所示):
consumer = KafkaConsumer(CONSUMER_TOPIC, group_id='ME2',
bootstrap_servers=[f"{KAFKA_SERVER_HOST}:{KAFKA_SERVER_PORT}"],
value_deserializer=lambda x: json.loads(x.decode('utf-8')),
enable_auto_commit=False,
auto_offset_reset='latest',
max_poll_records=1,
max_poll_interval_ms=300000)
with ThreadPoolExecutor(max_workers=10) as executor:
futures = []
for message in consumer:
futures.append(executor.submit(SOME_FUNCTION, ARG1, ARG2))
中间有一堆代码,但这段代码在这里并不重要,所以我跳过了。
现在,SOME_FUNCTION 来自另一个导入的 python 脚本(事实上,在后期阶段有一个导入层次结构)。重要的是,在这些脚本中的某个时刻,我调用了Multiprocessing 池,因为我需要对数据(SIMD - 单指令多数据)进行并行处理并使用 apply_async 函数来执行此操作。
for loop_message_chunk in loop_message_chunks:
res_list.append(self.pool.apply_async(self.one_matching.match, args=(hash_set, loop_message_chunk, fields)))
现在,我有 2 个版本的 runner/driver 程序:
-
基于Kafka(如上图)
- 此版本生成启动多处理的线程
听 Kafka -> 启动线程 -> 启动多处理
-
基于 REST(使用烧瓶通过 REST 调用实现相同的任务)
- 此版本不启动任何线程并立即调用多处理
监听 REST 端点 -> 启动多处理
您为什么要问 2 个跑步者/驱动程序脚本? - 这个微服务将被多个团队使用,有些团队想要基于同步 REST,而有些团队想要一个基于 KAFKA 的实时异步系统
当我从并行函数(上面示例中的self.one_matching.match)进行日志记录时,它在通过 REST 版本调用时有效,但在使用 KAFKA 版本调用时无效(基本上当多处理由线程启动时 - 它不起作用)。
还要注意,只有并行函数的日志记录不起作用。从运行器到调用 apply_async 的脚本的层次结构中的其余脚本(包括从线程内调用的脚本)成功记录。
其他细节:
- 我使用 yaml 文件配置记录器
- 我在运行脚本本身中为 KAFKA 或 REST 版本配置记录器
- 我在运行器脚本之后调用的所有其他脚本中执行
logging.getLogger,以让特定记录器记录到不同的文件
记录器配置(值替换为泛型,因为我无法获取确切名称):
version: 1
formatters:
simple:
format: '%(asctime)s | %(name)s | %(filename)s : %(funcName)s : %(lineno)d | %(levelname)s :: %(message)s'
custom1:
format: '%(asctime)s | %(filename)s :: %(message)s'
time-message:
format: '%(asctime)s | %(message)s'
handlers:
console:
class: logging.StreamHandler
level: DEBUG
formatter: simple
stream: ext://sys.stdout
handler1:
class: logging.handlers.TimedRotatingFileHandler
when: midnight
backupCount: 5
formatter: simple
level: DEBUG
filename: logs/logfile1.log
handler2:
class: logging.handlers.TimedRotatingFileHandler
when: midnight
backupCount: 30
formatter: custom1
level: INFO
filename: logs/logfile2.log
handler3:
class: logging.handlers.TimedRotatingFileHandler
when: midnight
backupCount: 30
formatter: time-message
level: DEBUG
filename: logs/logfile3.log
handler4:
class: logging.handlers.TimedRotatingFileHandler
when: midnight
backupCount: 30
formatter: time-message
level: DEBUG
filename: logs/logfile4.log
handler5:
class: logging.handlers.TimedRotatingFileHandler
when: midnight
backupCount: 5
formatter: simple
level: DEBUG
filename: logs/logfile5.log
loggers:
logger1:
level: DEBUG
handlers: [console, handler1]
propagate: no
logger2:
level: DEBUG
handlers: [console, handler5]
propagate: no
logger3:
level: INFO
handlers: [handler2]
propagate: no
logger4:
level: DEBUG
handlers: [console, handler3]
propagate: no
logger5:
level: DEBUG
handlers: [console, handler4]
propagate: no
kafka:
level: WARNING
handlers: [console]
propogate: no
root:
level: INFO
handlers: [console]
propogate: no
【问题讨论】:
-
我不知道我能回答为什么日志不能从一个线程启动的进程中工作,因为我希望它能够正常工作(大部分时间),并且然后有时会出现死锁(回复:6721)。我认为你可以摆脱线程但是使用aiokafka 在主(唯一)线程中创建一个 ProcessPoolExecutor,并根据需要从事件循环向它提交任务:docs.python.org/3/library/…
-
如果你想保持
SOME_FUNCTION不变(每次调用创建它自己的池而不是回调到全局ProcessPoolExecutor),它仍然应该以相同的方式工作。我只是在想,不继续创建和销毁单独的独立池可能会减少总开销。 -
似乎最简单的方法是使用带有logrotate的syslog,否则你需要在单独的进程中使用QueueListener和QueueHandler之类的东西,或者使用flask logger和你的kafka logger来登录不同的文件。
-
您难道不知道普通的日志记录不能很好地用于多处理吗?如果子进程是
forked,它可能会起作用,但如果它们是spawned,则不会。 QueueHandler 可能还不够,你需要 SocketHandler 来确定。您可以阅读此线程以了解更多信息stackoverflow.com/questions/64335940/…
标签: python multithreading logging multiprocessing python-logging