【问题标题】:Stateful processing in Apache Beam with key-value states具有键值状态的 Apache Beam 中的状态处理
【发布时间】:2021-08-13 16:14:35
【问题描述】:

我正在尝试使用 Apache Beam 实现一个有状态的进程。我浏览了 Kenneth Knowles 的两篇文章(Stateful processing with Apache BeamTimely (and Stateful) Processing with Apache Beam),但没有找到解决问题的方法。我正在使用 Python SDK。

特别是,我试图拥有一个包含键值对象的有状态 DoFn,我需要添加新元素,有时还需要删除一些元素。

我在 DoFn 课程中看到了 a solution may be to use a SetStateSpecTuple coder。问题是 SetSpaceSpec 没有类似“pop”的功能的选项。在我看来,删除元素的唯一方法是使用.clear() 将它们全部删除。 看来您不能只指定要使用此函数擦除的元素。

解决这个问题的一个机会可能是在我需要删除状态中的元素时清除和重写状态,但这对我来说似乎效率低下。

你知道如何有效地做到这一点吗?

Python 版本 3.8.7
apache-beam==2.29.0

【问题讨论】:

  • 例如,使用ReadModifyWriteStateSpec 并为 dict 自定义实现编码器怎么样?看看stackoverflow.com/a/45911664/8330018,它正在使用一个堆,但你可以从那里使用一些想法来实现.pop,例如
  • 嗨@TudorPlugaru,非常感谢您的回答。您对如何使用 Python SDK 类实现编码器有任何想法吗?您发布的示例使用 Java 的,我对它没有信心
  • 如果您需要代码示例,Apache Beam 存储库是一个不错的资源。关于自定义编码器的示例,您可以在这里找到它github.com/apache/beam/blob/master/sdks/python/apache_beam/…

标签: google-cloud-dataflow apache-beam stateful


【解决方案1】:

我按照@TudorPlugaru 的建议提出了这个建议。希望对其他人有用。

import json
from apache_beam.coders import Coder

class MyDictCoder(Coder):
    """ My custom dictionary coders """
    def encode(self, o):
        return json.dumps(o).encode()

    def decode(self, o):
        return json.loads(o.decode())

    def is_deterministic(self) -> bool:
        return True

在 DoFn 声明中

from apache_beam.transforms.userstate import ReadModifyWriteStateSpec

class MyDoFn(beam.DoFn):
    DICTSTATE= ReadModifyWriteStateSpec(name='dictstate', coder=MyDictCoder())
    
    def process(self, element, DictState=beam.DoFn.StateParam(DICTSTATE)):
        # Do something
        yield DictState

并在管道中添加这一行(如在Beam example 中所做的那样)

beam.coders.registry.register_coder(typing.Dict, MyDictCoder)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-25
    • 1970-01-01
    • 2021-03-08
    • 2012-07-31
    相关资源
    最近更新 更多