You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Apache Beam Python SDK的ParDo过滤PCollection元素?

Using ParDo to Filter Elements from a PCollection in Apache Beam

Absolutely, filtering elements with ParDo is a common task in Apache Beam, and it’s pretty straightforward once you see a concrete example. Let’s walk through how to do this, step by step.

Basic ParDo Filter Example

Say you have a PCollection of integers and want to keep only the even numbers. Here’s a complete implementation:

First, define a custom DoFn class that checks each element and keeps it only if it meets your condition:

import apache_beam as beam

class FilterEvenNumbers(beam.DoFn):
    def process(self, element):
        # Check if the element is even
        if element % 2 == 0:
            # Yield the element to include it in the output PCollection
            yield element

Then apply this DoFn to your PCollection using ParDo:

with beam.Pipeline() as pipeline:
    # Sample input data
    numbers = pipeline | beam.Create([1, 2, 3, 4, 5, 6, 7, 8])
    
    # Filter with ParDo
    even_numbers = numbers | beam.ParDo(FilterEvenNumbers())
    
    # Verify the result
    even_numbers | beam.Map(print)

How This Works

  • The FilterEvenNumbers class inherits from beam.DoFn, which is the base for all user-defined processing logic in ParDo.
  • The process method runs once per element. If the element passes your filter check (even number here), you yield it—this adds it to the output PCollection. Elements that don’t pass are ignored, so they get filtered out automatically.

Bonus: Built-in Filter Transform (For Simple Cases)

While ParDo gives you full control for complex logic, Beam has a built-in Filter transform that’s more concise for basic filtering. Here’s the same even-number example using it:

with beam.Pipeline() as pipeline:
    numbers = pipeline | beam.Create([1, 2, 3, 4, 5, 6, 7, 8])
    even_numbers = numbers | beam.Filter(lambda x: x % 2 == 0)
    even_numbers | beam.Map(print)

Use ParDo when you need to combine filtering with other logic (like modifying elements or emitting multiple outputs per input). For simple yes/no checks, Filter is the way to go.

Complex Filtering Example

If you need to filter based on dynamic or multi-part logic (e.g., filtering sentences that contain a specific keyword), here’s how to extend ParDo:

class FilterContainsKeyword(beam.DoFn):
    def __init__(self, keyword):
        self.keyword = keyword
    
    def process(self, element):
        # Case-insensitive check for the keyword
        if self.keyword.lower() in element.lower():
            yield element

with beam.Pipeline() as pipeline:
    sentences = pipeline | beam.Create([
        "Apache Beam is a unified data processing framework",
        "ParDo is a core transform in Beam",
        "Filter is great for simple cases",
        "Python is a supported SDK for Beam"
    ])
    
    filtered_sentences = sentences | beam.ParDo(FilterContainsKeyword("beam"))
    filtered_sentences | beam.Map(print)

This will output all sentences that include the word "beam", regardless of capitalization.


内容的提问来源于stack exchange,提问作者Pablo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 06:51:10