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

PyFlink中使用SourceFunction实现数据生成器的技术问询

问题原因解释

你遇到的 AttributeError: 'DataGenerator' object has no attribute '_get_object_id' 错误,是因为PyFlink的Python API无法直接继承Java端的SourceFunction接口。这是由于PyFlink的Python-JVM桥接机制要求对象具备内部通信所需的方法,而自定义Python类无法满足这一要求。

可行的OOP实现方式

在新版PyFlink(1.16+)中,推荐使用Python原生的Source抽象类来实现自定义无界数据源,完全不需要依赖Java或Jython。以下是类似Java版TaxiRideGenerator的OOP风格实现示例:

完整示例代码

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.source import (
    Source, SourceReader, SplitEnumerator,
    SingleSplitSourceReader, SingleSplitEnumerator
)
from pyflink.datastream.source.split import SingleSplit
from typing import Iterator, Optional
import time
import random
from dataclasses import dataclass

# 定义TaxiRide数据结构
@dataclass
class TaxiRide:
    ride_id: int
    is_start: bool
    passenger_count: int
    timestamp: int

# 自定义数据源类,实现Source接口
class TaxiRideGeneratorSource(Source[TaxiRide]):
    def __init__(self, max_events: Optional[int] = None):
        self.max_events = max_events  # 可选:限制生成的事件总数

    def create_reader(self, reader_context) -> SourceReader[TaxiRide, SingleSplit]:
        return TaxiRideGeneratorReader(self.max_events)

    def create_enumerator(self, enumerator_context) -> SplitEnumerator[SingleSplit]:
        # 使用SingleSplitEnumerator实现非并行数据源
        return SingleSplitEnumerator(enumerator_context, SingleSplit("taxi-ride-split"))

    def restore_enumerator(self, enumerator_context, checkpoint) -> SplitEnumerator[SingleSplit]:
        # 恢复枚举器(用于容错)
        return self.create_enumerator(enumerator_context)

# 自定义数据源读取器,负责生成数据
class TaxiRideGeneratorReader(SingleSplitSourceReader[TaxiRide]):
    def __init__(self, max_events: Optional[int]):
        self.max_events = max_events
        self.event_count = 0
        self.running = True

    def poll_next(self) -> Iterator[TaxiRide]:
        # 停止条件:达到最大事件数或被关闭
        if (self.max_events is not None and self.event_count >= self.max_events) or not self.running:
            return
        
        # 生成随机出租车行程数据
        ride_id = random.randint(1, 10000)
        is_start = random.choice([True, False])
        passenger_count = random.randint(1, 4)
        timestamp = int(time.time() * 1000)  # 毫秒级时间戳
        
        self.event_count += 1
        yield TaxiRide(ride_id, is_start, passenger_count, timestamp)
        
        # 模拟事件间隔,避免生成过快
        time.sleep(0.1)

    def is_running(self) -> bool:
        return self.running

    def close(self):
        self.running = False

if __name__ == "__main__":
    # 初始化执行环境
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(1)  # 非并行数据源设置为1

    # 创建自定义数据源
    taxi_ride_source = TaxiRideGeneratorSource()

    # 将数据源添加到数据流
    ds = env.from_source(
        source=taxi_ride_source,
        watermark_strategy=None,  # 可根据需求添加水位线策略
        source_name="TaxiRideGenerator"
    )

    # 打印输出结果
    ds.print()

    # 执行作业
    env.execute("Unbounded Taxi Ride Generator Job")

关键说明

  1. 数据结构定义:使用dataclasses定义TaxiRide类,PyFlink可自动识别并序列化该类型。
  2. Source接口实现:TaxiRideGeneratorSource负责创建读取器和枚举器,SingleSplitEnumerator适用于无需并行拆分的数据源。
  3. 数据生成逻辑:TaxiRideGeneratorReader的poll_next方法负责生成并返回数据,is_running控制数据源的运行状态。
  4. 无界/有界切换:通过max_events参数可轻松将无界数据源转为有界数据源。

关于Source/SinkFunction接口的说明

PyFlink的Python API不支持直接使用Java端的SourceFunction/SinkFunction接口。如需自定义Sink,可参考Python原生的Sink抽象类,实现方式与上述Source类似。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 16:15:54