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

PyFlink是否支持max或max_by函数?Java实现滚动取最大值后Python版本报错求助

解决PyFlink滚动取最大值的问题 & API支持说明

首先明确:PyFlink的DataStream API确实没有像Java Flink那样给KeyedStream提供直接的max()或max_by()方法,这就是你遇到AttributeError的原因。不过我们有两种可靠的方式实现和你Java代码完全一致的滚动取最大值功能,下面分别说明:


方案1:使用DataStream API自定义AggregateFunction

我们可以通过实现AggregateFunction来模拟Java中max()的行为——按传感器ID分组后,滚动保留当前温度的最大值,同时保留该最大值第一次出现时的timestamp(和你Java代码的输出逻辑完全对齐)。

完整代码示例

from pyflink.datastream.functions import AggregateFunction
from pyflink.datastream import StreamExecutionEnvironment, RuntimeExecutionMode

class RollingMaxTemperature(AggregateFunction):
    # 定义累加器,保存当前分组的最大温度、对应ID和timestamp
    class Accumulator:
        def __init__(self):
            self.id = None
            self.max_temp = -float('inf')
            self.timestamp = None

    def create_accumulator(self):
        return self.Accumulator()

    def add(self, value, accumulator):
        # 将输入的温度转为float类型用于比较
        current_temp = value["temperature"]
        # 当当前温度大于已记录的最大值时,更新状态
        if current_temp > accumulator.max_temp:
            accumulator.id = value["id"]
            accumulator.max_temp = current_temp
            accumulator.timestamp = value["timestamp"]
        # 温度相等时保留原有记录,和Java的max行为一致

    def get_result(self, accumulator):
        # 返回最终的聚合结果
        return {
            "id": accumulator.id,
            "timestamp": accumulator.timestamp,
            "temperature": accumulator.max_temp
        }

    def merge(self, accumulator1, accumulator2):
        # 并行场景下合并两个累加器的状态,取温度更大的那个;温度相同时保留第一个的timestamp
        if accumulator1.max_temp > accumulator2.max_temp:
            return accumulator1
        elif accumulator2.max_temp > accumulator1.max_temp:
            return accumulator2
        else:
            return accumulator1

def f_map():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
    env.set_parallelism(1)
    
    # 读取数据源并转换为结构化字典
    ds = env.read_text_file("/Users/yinbaidong/PycharmProjects/pythonProject2/wudu.txt")
    new_ds = ds.map(
        lambda i: {
            "id": i.split(",")[0],
            "timestamp": i.split(",")[1],
            "temperature": float(i.split(",")[2])
        }
    )
    
    # 分组后应用自定义聚合函数
    result_ds = new_ds.key_by(lambda i: i.get("id")).aggregate(RollingMaxTemperature())
    
    result_ds.print()
    env.execute("Rolling Max Temperature Job")

if __name__ == '__main__':
    f_map()

方案2:使用Table API(更简洁)

PyFlink的Table API提供了和Java Flink类似的max()和max_by()函数,实现起来更直观,代码量更少,推荐优先使用这种方式。

完整代码示例

from pyflink.table import StreamTableEnvironment, EnvironmentSettings
from pyflink.datastream import StreamExecutionEnvironment

def table_api_approach():
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_runtime_mode(RuntimeExecutionMode.STREAMING)
    env.set_parallelism(1)
    
    # 创建流式Table环境
    settings = EnvironmentSettings.new_instance().in_streaming_mode().build()
    t_env = StreamTableEnvironment.create(env, environment_settings=settings)
    
    # 注册文件数据源为临时表
    t_env.execute_sql("""
        CREATE TABLE sensor_data (
            id STRING,
            timestamp BIGINT,
            temperature DOUBLE
        ) WITH (
            'connector' = 'filesystem',
            'path' = '/Users/yinbaidong/PycharmProjects/pythonProject2/wudu.txt',
            'format' = 'csv'
        )
    """)
    
    # 按ID分组,滚动取温度最大值,并保留该最大值第一次出现的timestamp(和Java逻辑对齐)
    result_table = t_env.from_path("sensor_data") \
        .group_by("id") \
        .select(
            "id",
            max("temperature").alias("temperature"),
            first_value("timestamp").filter_(temperature == max(temperature)).alias("timestamp")
        )
    
    # 转换为DataStream并打印输出
    result_ds = t_env.to_data_stream(result_table)
    result_ds.print()
    
    env.execute("Rolling Max with Table API")

if __name__ == '__main__':
    table_api_approach()

验证结果

两种方案的输出都会和你Java代码的结果一致,比如:

+I[sensor_1, 1547718199, 35.8]
+I[sensor_2, 1547718299, 35.8]
-U[sensor_2, 1547718299, 35.8]
+I[sensor_2, 1547718299, 35.8]  # 温度未超过当前最大值,无变化
-U[sensor_2, 1547718299, 35.8]
+I[sensor_2, 1547718599, 36.8]  # 温度更新为更大值
-U[sensor_1, 1547718199, 35.8]
+I[sensor_1, 1547718699, 41.8]
-U[sensor_1, 1547718699, 41.8]
+I[sensor_1, 1547718699, 41.8]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 00:32:37