// Flink Java 示例代码:Word Count
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;
public class WordCount {
public static void main(String[] args) throws Exception {
// 设置执行环境
final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
// 从文本文件中读取数据
DataSet<String> text = env.readTextFile("file:///path/to/input");
// 执行 word count 操作
DataSet<Tuple2<String, Integer>> counts = text
.flatMap(new Tokenizer())
.groupBy(0)
.sum(1);
// 将结果输出到文件
counts.writeAsCsv("file:///path/to/output", "\n", " ");
// 执行程序
env.execute("Word Count Example");
}
// 自定义的 FlatMap 函数,用于将每行文本拆分为单词
public static final class Tokenizer implements FlatMapFunction<String, Tuple2<String, Integer>> {
@Override
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) {
// 将每一行按空格分割成单词
String[] tokens = value.toLowerCase().split("\\W+");
for (String token : tokens) {
if (token.length() > 0) {
out.collect(new Tuple2<>(token, 1));
}
}
}
}
}
ExecutionEnvironment:
getExecutionEnvironment() 方法会根据上下文返回适当的执行环境(例如本地环境或集群环境)。DataSet:
DataSet 是 Flink 中的核心抽象之一,表示分布式的数据集。这里我们使用 readTextFile 方法从指定路径读取文本文件,并将其转换为 DataSet<String>。FlatMapFunction:
FlatMapFunction 是一个接口,用于将输入元素转换为零个、一个或多个输出元素。在这个例子中,我们将每行文本拆分为多个单词,并为每个单词生成一个 (word, 1) 的元组。groupBy 和 sum:
groupBy(0) 表示按照元组的第一个字段(即单词)进行分组,sum(1) 则对每个分组中的第二个字段(即计数)进行求和。writeAsCsv:
\n 和空格。env.execute:
execute 方法来启动 Flink 程序的执行,并给定一个作业名称 "Word Count Example"。这个示例展示了如何使用 Apache Flink 进行简单的词频统计(Word Count)。
上一篇:java 字符串转对象
Laravel PHP 深圳智简公司。版权所有©2023-2043 LaravelPHP 粤ICP备2021048745号-3
Laravel 中文站