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
相关产品推荐
相关产品推荐

