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

Apache Beam流式Pipeline转DataFrame报错:RuntimeError: NotImplementedError

问题描述

我正处于Apache Beam学习阶段,编写了一个从Pub/Sub读取数据并写入BigQuery的流式Pipeline,但在将PCollection转换为DataFrame以进行后续数据转换时,遇到错误:

RuntimeError: NotImplementedError [while running 'BatchElements(messages)']

请问我遗漏了什么?

我的代码:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.io.gcp.bigquery import WriteToBigQuery
from apache_beam.dataframe.convert import to_dataframe, to_pcollection
import os
import typing

os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "my_json_file.json"


# Do function for passing in pubsub message from the subscription
# Define a NamedTuple schema
class BmsSchema(typing.NamedTuple):
    ident: str


class ParsePubSubMessage(beam.DoFn):
    def process(self, message):
        import json
        # Creating the main_dict that has all the columns
        all_columns = ['ident']
        main_dict = dict(zip(all_columns, [None] * len(all_columns)))
        # Parse the JSON message
        record = json.loads(message.decode('utf-8'))
        main_dict.update(record)

        yield {
            'ident': main_dict["ident"]
        }


def run():
    # Define pipeline options
    options = PipelineOptions(
        project='dwingestion',
        runner='DirectRunner',
        streaming=True,  # Enable streaming mode
        temp_location='gs://........./temp',
        staging_location='gs://....../staging',
        region='europe-west1',
        job_name='flesp-streaming-pipeline-dataflow-test'
    )

    # Set streaming mode
    options.view_as(StandardOptions).streaming = True

    # Pub/Sub subscription
    input_subscription = 'projects/...../subscriptions/flespi_data_streaming'

    table_schema = {
        "fields": [
            {"name": "ident", "type": "STRING", "mode": "NULLABLE"}
        ]
    }

    # Create the pipeline
    with beam.Pipeline(options=options) as p:
        # Read from Pub/Sub and parse the messages
        messages = (p
                    | 'Read from PubSub' >> beam.io.ReadFromPubSub(subscription=input_subscription)
                    | 'Parse PubSub Message' >> beam.ParDo(ParsePubSubMessage())
                    | 'Attaching the schema' >> beam.Map(lambda x: BmsSchema(**x)).with_output_types(BmsSchema)
                    )

        # Convert the messages to df
        df = to_dataframe(messages)
        transformed_pcol = to_pcollection(df)
        # Write to BigQuery with schema autodetect
        transformed_pcol | 'Write to BigQuery' >> WriteToBigQuery(
            table='.........flesp_table_test_4',
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
            schema=table_schema,
            custom_gcs_temp_location='gs://........../temp'
        )


if __name__ == '__main__':
    run()
问题原因与解决方案

核心原因

你遇到的错误是因为Apache Beam的DataFrame转换工具(to_dataframe/to_pcollection)不支持无界流式PCollection。to_dataframe内部会执行BatchElements操作来将数据打包成批,但这个操作在流式模式下并未实现——流式数据是无限的,无法直接适配DataFrame需要固定批次的结构特性。

解决办法

方案1:用Beam原生操作替代DataFrame转换(推荐)

这是流式Pipeline的标准实现方式,完全规避DataFrame的限制,直接对PCollection进行处理:

import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.io.gcp.bigquery import WriteToBigQuery
import os
import typing

os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "my_json_file.json"


class BmsSchema(typing.NamedTuple):
    ident: str


class ParsePubSubMessage(beam.DoFn):
    def process(self, message):
        import json
        record = json.loads(message.decode('utf-8'))
        # 直接生成符合Schema的NamedTuple
        yield BmsSchema(ident=record.get("ident"))


def run():
    options = PipelineOptions(
        project='dwingestion',
        runner='DirectRunner',
        streaming=True,
        temp_location='gs://........./temp',
        staging_location='gs://....../staging',
        region='europe-west1',
        job_name='flesp-streaming-pipeline-dataflow-test'
    )
    options.view_as(StandardOptions).streaming = True

    input_subscription = 'projects/...../subscriptions/flespi_data_streaming'

    table_schema = {
        "fields": [
            {"name": "ident", "type": "STRING", "mode": "NULLABLE"}
        ]
    }

    with beam.Pipeline(options=options) as p:
        messages = (p
                    | 'Read from PubSub' >> beam.io.ReadFromPubSub(subscription=input_subscription)
                    | 'Parse PubSub Message' >> beam.ParDo(ParsePubSubMessage())
                    )

        # 直接将PCollection写入BigQuery,无需DataFrame转换
        messages | 'Write to BigQuery' >> WriteToBigQuery(
            table='.........flesp_table_test_4',
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
            create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
            schema=table_schema,
            custom_gcs_temp_location='gs://........../temp'
        )


if __name__ == '__main__':
    run()

方案2:给流式数据加窗口,转为有界批次(业务允许时使用)

如果必须使用DataFrame处理逻辑,可以通过**窗口(Window)**将无限流式数据分割为固定时间窗口的有界批次,再进行DataFrame转换。注意这会引入数据延迟,需根据业务需求调整窗口参数:

# 在转换为DataFrame前添加窗口操作
windowed_messages = messages | '60s Fixed Window' >> beam.WindowInto(beam.window.FixedWindows(60))

df = to_dataframe(windowed_messages)
# 在这里添加DataFrame转换逻辑(如筛选、计算等)
transformed_pcol = to_pcollection(df)

该方式本质是将流式数据拆分为多个小批处理任务,会改变Pipeline的语义(如延迟、数据处理粒度),需谨慎评估。


内容的提问来源于stack exchange,提问作者oyugi.collins

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 19:30:08