Hacktoberfest 2026:维护者为十月标记出来的 issue,仍然开放、适合新手。 浏览 Hacktoberfest issue

[Bug]: [RunInference] max_models_per_worker_hint is not enforced

未关闭
#40,468 1 条评论 0 个 reaction 已指派 1 人 在 GitHub 查看

维护者通常 1 天内回复

@schizophrenicmaniac 已经在做这个了。

开始于 2026年10月10日。

  • #40500 来自 @schizophrenicmaniac —— 未关闭

评估

难度
4/5
预计耗时
3-5 天
新手友好度
48/100
Issue 类型
缺陷
描述清晰度
基本清楚
活跃度
活跃
技术栈
python
领域
backend

调研方向

Start in sdks/python/apache_beam/ml/inference/base.py around the max-workers logic at lines 852–857, and trace how locks and deserialized RunInference handlers interact with the shared _ModelHandlerManager. Use the issue’s reproduction as a starting point and check existing Python SDK inference tests. Done when the model limit remains at the configured hint across multiple handler copies and the regression is covered by a test.

由索引模型根据 Issue 内容生成。

描述

bug P2 python
What happened?

The intent of max_models_per_worker hint is to limit the number of models that will be loaded per SDK process.

The logic to increment allowed max number of workers https://github.com/apache/beam/blob/dabcf50ffbf532bb30a1233927cd1edb6ae067bb/sdks/python/apache_beam/ml/inference/base.py#L852-L857 appears to be flawed since:

  1. The lock acquisition always succeeds (it's a new lock instance)
  2. If there is more than 1 process bundle descriptor over the life time of the SDK process, we might have more than 1 copy of the deserialized RunInference DoFn with unpickled ModelHandler instance, which won't persist self._max_models_per_worker_hint = None from a prior initialization:

AI repro:

from apache_beam.internal import pickler
from apache_beam.ml.inference import base
class Model:
  def predict(self, x):
    return x
class Handler(base.ModelHandler):
  def load_model(self):
    return Model()
  def run_inference(self, batch, model, inference_args=None):
    return [model.predict(x) for x in batch]
mhs = [base.KeyModelMapping([k], Handler()) for k in ('a', 'b', 'c')]
keyed_handler = base.KeyedModelHandler(mhs, max_models_per_worker_hint=1)
# RunInference shares one _ModelHandlerManager per transform across all
# DoFn instances (and processes) via MultiProcessShared.
manager = keyed_handler.load_model()
# 5 DoFn instances in one process (harness threads, re-created bundle
# processors); each deserializes its own copy of the model handler.
for _ in range(5):
  handler_copy = pickler.roundtrip(keyed_handler)
  handler_copy.override_metrics('ns')
  handler_copy.run_inference([('a', 1), ('b', 2), ('c', 3)], manager)
print('model limit:', manager._max_models)  # 5, expected 1
print('models in memory:', len(manager._tag_map))  # 3

Issue Priority

Priority: 2 (default / most bugs should be filed as P2)

Issue Components
  • Component: Python SDK
  • Component: Java SDK
  • Component: Go SDK
  • Component: Typescript SDK
  • Component: IO connector
  • Component: Beam YAML
  • Component: Beam examples
  • Component: Beam playground
  • Component: Beam katas
  • Component: Website
  • Component: Infrastructure
  • Component: Spark Runner
  • Component: Flink Runner
  • Component: Prism Runner
  • Component: Twister2 Runner
  • Component: Hazelcast Jet Runner
  • Component: Google Cloud Dataflow Runner
主要语言
Java
星标
8.7k
派生
4.7k
平均合并
2 天 7 小时
30 天内合并 PR
242

环境准备

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

apache/beam 的其他 Issue

查看 apache/beam 的全部 Issue

相似的 Issue

更多 Java Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。