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

如何使用Flink ML迭代机制?原DataStream迭代示例迁移疑问

以下是一个简单的非机器学习迭代示例,实现和旧IterativeStream类似的元素递减逻辑:

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.ml.common.datastream.DataStreamList;
import org.apache.flink.ml.common.iteration.IterationBody;
import org.apache.flink.ml.common.iteration.Iterations;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.FilterFunction;

public class NonMLIterationDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // 初始变量流:输入0-5的整数
        DataStream<Long> initVars = env.generateSequence(0, 5);

        // 定义迭代逻辑体
        IterationBody iterationBody = (variableStreams, dataStreams) -> {
            // 获取当前迭代的变量流
            DataStream<Long> currentStream = variableStreams.get(0);
            
            // 每个元素减1
            DataStream<Long> minusOne = currentStream.map((MapFunction<Long, Long>) value -> value - 1);
            
            // 筛选需要继续迭代的元素(>0)
            DataStream<Long> toContinue = minusOne.filter((FilterFunction<Long>) value -> value > 0);
            
            // 筛选需要输出的元素(<=0)
            DataStream<Long> toOutput = minusOne.filter((FilterFunction<Long>) value -> value <= 0);

            // 返回迭代结果:继续迭代的流 + 输出流
            return Iterations.Result.of(DataStreamList.of(toContinue), DataStreamList.of(toOutput));
        };

        // 执行无界迭代,设置最大迭代次数防止无限循环
        DataStreamList outputStreams = Iterations.iterateUnboundedStreams(
                DataStreamList.of(initVars),
                DataStreamList.empty(), // 无外部数据流,传空列表
                iterationBody,
                100
        );

        // 打印输出结果
        outputStreams.get(0).print("Final Output");

        env.execute("Non-ML Iteration Example");
    }
}

2. 复现旧IterativeStream示例的方案

你提供的旧IterativeStream示例完全可以用Flink ML的Iterations.iterateUnboundedStreams复现,以下是对应实现及参数说明:

复现代码

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.ml.common.datastream.DataStreamList;
import org.apache.flink.ml.common.iteration.IterationBody;
import org.apache.flink.ml.common.iteration.Iterations;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.FilterFunction;

public class IterativeStreamMigrationDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // 初始变量流:对应旧示例中的someIntegers
        DataStream<Long> someIntegers = env.generateSequence(0, 1000);

        // 定义迭代逻辑体
        IterationBody iterationBody = (variableStreams, dataStreams) -> {
            // 获取当前迭代的输入流(上一轮迭代返回的stillGreaterThanZero)
            DataStream<Long> iterationInput = variableStreams.get(0);
            
            // 对应旧示例中的minusOne
            DataStream<Long> minusOne = iterationInput.map(new MapFunction<Long, Long>() {
                @Override
                public Long map(Long value) throws Exception {
                    return value - 1;
                }
            });

            // 对应旧示例中的stillGreaterThanZero:继续迭代的流
            DataStream<Long> stillGreaterThanZero = minusOne.filter(new FilterFunction<Long>() {
                @Override
                public boolean filter(Long value) throws Exception {
                    return value > 0;
                }
            });

            // 对应旧示例中的lessThanZero:最终输出的流
            DataStream<Long> lessThanZero = minusOne.filter(new FilterFunction<Long>() {
                @Override
                public boolean filter(Long value) throws Exception {
                    return value <= 0;
                }
            });

            // 返回迭代结果:继续迭代的变量流 + 输出流
            return Iterations.Result.of(DataStreamList.of(stillGreaterThanZero), DataStreamList.of(lessThanZero));
        };

        // 执行迭代
        DataStreamList outputStreams = Iterations.iterateUnboundedStreams(
                DataStreamList.of(someIntegers), // initVariableStreams:初始输入流
                DataStreamList.empty(), // dataStreams:旧示例无外部数据流,传空列表
                iterationBody,
                1000 // 最大迭代次数,避免极端情况无限循环
        );

        // 打印输出结果
        outputStreams.get(0).print("Less Than Zero");

        env.execute("Migrate IterativeStream to Flink ML Iteration");
    }
}

参数对应关系说明

  • initVariableStreams:对应旧示例中someIntegers.iterate()的初始流,是迭代的起始输入变量。
  • dataStreams:用于传递每次迭代都需要处理的外部静态数据流(比如固定的参考数据),你的旧示例没有这类外部流,所以传DataStreamList.empty()即可。
  • IterationBody:封装每次迭代的核心逻辑,接收当前的变量流和外部数据流,处理后返回两部分:
    • newVariableStreams:对应旧示例中iteration.closeWith(stillGreaterThanZero)的流,即需要回到迭代继续处理的部分。
    • outputStreams:对应旧示例中lessThanZero的流,即迭代终止后输出的结果。

内容的提问来源于stack exchange,提问作者K.M

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 18:59:49