如何使用Flink ML迭代机制?原DataStream迭代示例迁移疑问
Flink ML迭代替代DataStream IterativeStream问题解答
1. 非ML场景的Flink ML迭代可运行示例
以下是一个简单的非机器学习迭代示例,实现和旧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
相关产品推荐
相关产品推荐

