使用apache_beam.utils.shared创建查找数据时遇类型错误
问题原因
Apache Beam的shared.Shared模块内部依赖Python的weakref机制管理共享对象,但Python的内置基础类型(比如list、dict、int等)不支持创建弱引用。官方示例直接返回list给acquire方法,这在实际运行时会触发类型错误。
解决方案
把需要共享的list包装到一个自定义类的实例中——自定义类的对象支持弱引用,可以被Shared模块正常处理。修改后的代码如下:
import apache_beam as beam from apache_beam.utils import shared from log_elements import LogElements # 自定义包装类,用于持有list数据,使其支持弱引用 class ListHolder: def __init__(self, data): self.data = data class GetNthStringFn(beam.DoFn): def __init__(self, shared_handle): self._shared_handle = shared_handle def process(self, element): def initialize_list(): # 将生成的list包装到ListHolder实例中返回 return ListHolder([str(i) for i in range(1000000)]) giant_list_holder = self._shared_handle.acquire(initialize_list) # 通过实例的data属性访问原list yield giant_list_holder.data[element] with beam.Pipeline() as p: shared_handle = shared.Shared() (p | beam.Create([2, 4, 6, 8]) | beam.ParDo(GetNthStringFn(shared_handle)) | LogElements())
修改说明
- 新增
ListHolder类:仅作为容器存储list数据,利用自定义类实例可被弱引用的特性绕过限制 - 调整
initialize_list的返回值:不再直接返回list,而是返回持有该list的ListHolder实例 - 访问共享数据时:通过
ListHolder实例的data属性获取原list
内容的提问来源于stack exchange,提问作者mayur pokiya
相关产品推荐
相关产品推荐

