- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我在 Flink (Java) 中有这个程序,它可以计算数据流中的不同单词。我使用计数单词的示例来实现,并在同一时间应用另一个窗口来评估不同的值。该程序运行良好。但是,我担心我正在使用两个窗口来处理不同的计数。第一个窗口计算单词数,第二个窗口将单词数切换为 1
,并将单词切换为 Tuple2
的第二个元素。我数了数 key 的数量。这是我的程序的输入和输出:
// input:
aaa
aaa
bbb
ccc
bbb
aaa
output:
(3,bbb-ccc-aaa)
如果我删除第二个窗口,它会显示每个键的所有评估并保存前一个窗口的状态。
// input:
aaa
aaa
bbb
ccc
bbb
aaa
// output:
3> (1,bbb)
3> (2,bbb-aaa)
3> (3,bbb-aaa-ccc)
// wait the first window to be evaluated.
// input:
aaa
aaa
bbb
ccc
bbb
aaa
// output:
3> (4,bbb-aaa-ccc-ccc)
3> (5,bbb-aaa-ccc-ccc-bbb)
3> (6,bbb-aaa-ccc-ccc-bbb-aaa)
我的程序:
public class WordCountDistinctSocketFilterQEP {
public WordCountDistinctSocketFilterQEP() throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// @formatter:off
env.socketTextStream("localhost", 9000)
.flatMap(new SplitterFlatMap())
.keyBy(new MyKeySelector())
.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
.reduce(new CountReduceFunction())
.map(new SwapMapFunction())
.keyBy(0)
.window(TumblingProcessingTimeWindows.of(Time.seconds(5))) // TESTING REMOVING THIS WINDOW
.reduce(new CountDistinctFunction())
.print();
// @formatter:on
String executionPlan = env.getExecutionPlan();
System.out.println("ExecutionPlan ........................ ");
System.out.println(executionPlan);
System.out.println("........................ ");
// dataStream.print();
env.execute("WordCountDistinctSocketFilterQEP");
}
public static class SwapMapFunction implements MapFunction<Tuple2<String, Integer>, Tuple2<Integer, String>> {
private static final long serialVersionUID = 5148172163266330182L;
@Override
public Tuple2<Integer, String> map(Tuple2<String, Integer> value) throws Exception {
return Tuple2.of(1, value.f0);
}
}
public static class SplitterFlatMap implements FlatMapFunction<String, Tuple2<String, Integer>> {
private static final long serialVersionUID = 3121588720675797629L;
@Override
public void flatMap(String sentence, Collector<Tuple2<String, Integer>> out) throws Exception {
for (String word : sentence.split(" ")) {
out.collect(new Tuple2<String, Integer>(word, 1));
}
}
}
public static class MyKeySelector implements KeySelector<Tuple2<String, Integer>, String> {
private static final long serialVersionUID = 2787589690596587044L;
@Override
public String getKey(Tuple2<String, Integer> value) throws Exception {
return value.f0;
}
}
public static class CountReduceFunction implements ReduceFunction<Tuple2<String, Integer>> {
private static final long serialVersionUID = 8541031982462158730L;
@Override
public Tuple2<String, Integer> reduce(Tuple2<String, Integer> value1, Tuple2<String, Integer> value2)
throws Exception {
return Tuple2.of(value1.f0, value1.f1 + value2.f1);
}
}
public static class CountDistinctFunction implements ReduceFunction<Tuple2<Integer, String>> {
private static final long serialVersionUID = -7077952757215699563L;
@Override
public Tuple2<Integer, String> reduce(Tuple2<Integer, String> value1, Tuple2<Integer, String> value2)
throws Exception {
return Tuple2.of(value1.f0 + value2.f0, value1.f1 + "-" + value2.f1);
}
}
}
最佳答案
ReduceFunctions
与 Collections
更好地合作( Maps
、 Lists
、 Sets
)。如果将每个单词映射到一个元素 Set
,你可以写一个 ReduceFunction
运行于 Set<String>
然后你可以用一个 ReduceFunction
来做到这一点而不是两个。
所以有splitterFlatMap
返回一系列由一个元素组成的长 Set<String>
, MyKeySelector
返回每个集合的第一个元素。窗口函数很好,更改reduce函数以匹配Set<String>
类型,函数的核心是 value1.addAll(value2)
。此时,您已经获得了输入中所有唯一单词的集合,这些单词分布在您正在运行的多个并行任务中。根据完成后您将所有这些数据放在哪里,这可能就足够了。否则,您可以在其末尾放置一个全局窗口,并再次使用相同的reduce函数(解释如下)
你的第二个问题是这不会按原样扩展。在某种程度上,这更像是一个哲学问题。如果不让每个并行实例都与其他实例通信,您就无法真正获得跨并行实例的全局计数。不过,您可以做的是通过实际单词对拆分单词流进行键控,然后使用(并行)键控、窗口 ReduceFunction
获取每个键组中不同单词的列表。然后你可以再吃一个ReduceFunction
这不是并行的,它结合了并行结果的结果。您还希望第二个窗口也打开; WindowFunctions
在触发之前等待所有上游运算符达到正确的水印,因此窗口将确保您的非并行运算符接收来自每个并行运算符的输入。非并行运算符上的聚合是简单的串联,因为一开始的键控保证给定的单词将恰好存在于一个并行槽中。
很明显,单个非并行运算符可能会出现瓶颈,但负载规模与不同单词的总数有关,实际上,由于英语的工作方式,负载规模可能仅限于 10k 单词左右.
关于java - 如何提高 Flink 中数据流实现的不同计数?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/56524962/
背景: 我最近一直在使用 JPA,我为相当大的关系数据库项目生成持久层的轻松程度给我留下了深刻的印象。 我们公司使用大量非 SQL 数据库,特别是面向列的数据库。我对可能对这些数据库使用 JPA 有一
我已经在我的 maven pom 中添加了这些构建配置,因为我希望将 Apache Solr 依赖项与 Jar 捆绑在一起。否则我得到了 SolarServerException: ClassNotF
interface ITurtle { void Fight(); void EatPizza(); } interface ILeonardo : ITurtle {
我希望可用于 Java 的对象/关系映射 (ORM) 工具之一能够满足这些要求: 使用 JPA 或 native SQL 查询获取大量行并将其作为实体对象返回。 允许在行(实体)中进行迭代,并在对当前
好像没有,因为我有实现From for 的代码, 我可以转换 A到 B与 .into() , 但同样的事情不适用于 Vec .into()一个Vec . 要么我搞砸了阻止实现派生的事情,要么这不应该发
在 C# 中,如果 A 实现 IX 并且 B 继承自 A ,是否必然遵循 B 实现 IX?如果是,是因为 LSP 吗?之间有什么区别吗: 1. Interface IX; Class A : IX;
就目前而言,这个问题不适合我们的问答形式。我们希望答案得到事实、引用资料或专业知识的支持,但这个问题可能会引发辩论、争论、投票或扩展讨论。如果您觉得这个问题可以改进并可能重新打开,visit the
我正在阅读标准haskell库的(^)的实现代码: (^) :: (Num a, Integral b) => a -> b -> a x0 ^ y0 | y0 a -> b ->a expo x0
我将把国际象棋游戏表示为 C++ 结构。我认为,最好的选择是树结构(因为在每个深度我们都有几个可能的移动)。 这是一个好的方法吗? struct TreeElement{ SomeMoveType
我正在为用户名数据库实现字符串匹配算法。我的方法采用现有的用户名数据库和用户想要的新用户名,然后检查用户名是否已被占用。如果采用该方法,则该方法应该返回带有数据库中未采用的数字的用户名。 例子: “贾
我正在尝试实现 Breadth-first search algorithm , 为了找到两个顶点之间的最短距离。我开发了一个 Queue 对象来保存和检索对象,并且我有一个二维数组来保存两个给定顶点
我目前正在 ika 中开发我的 Python 游戏,它使用 python 2.5 我决定为 AI 使用 A* 寻路。然而,我发现它对我的需要来说太慢了(3-4 个敌人可能会落后于游戏,但我想供应 4-
我正在寻找 Kademlia 的开源实现C/C++ 中的分布式哈希表。它必须是轻量级和跨平台的(win/linux/mac)。 它必须能够将信息发布到 DHT 并检索它。 最佳答案 OpenDHT是
我在一本书中读到这一行:-“当我们要求 C++ 实现运行程序时,它会通过调用此函数来实现。” 而且我想知道“C++ 实现”是什么意思或具体是什么。帮忙!? 最佳答案 “C++ 实现”是指编译器加上链接
我正在尝试使用分支定界的 C++ 实现这个背包问题。此网站上有一个 Java 版本:Implementing branch and bound for knapsack 我试图让我的 C++ 版本打印
在很多情况下,我需要在 C# 中访问合适的哈希算法,从重写 GetHashCode 到对数据执行快速比较/查找。 我发现 FNV 哈希是一种非常简单/好/快速的哈希算法。但是,我从未见过 C# 实现的
目录 LRU缓存替换策略 核心思想 不适用场景 算法基本实现 算法优化
1. 绪论 在前面文章中提到 空间直角坐标系相互转换 ,测绘坐标转换时,一般涉及到的情况是:两个直角坐标系的小角度转换。这个就是我们经常在测绘数据处理中,WGS-84坐标系、54北京坐标系
在软件开发过程中,有时候我们需要定时地检查数据库中的数据,并在发现新增数据时触发一个动作。为了实现这个需求,我们在 .Net 7 下进行一次简单的演示. PeriodicTimer .
二分查找 二分查找算法,说白了就是在有序的数组里面给予一个存在数组里面的值key,然后将其先和数组中间的比较,如果key大于中间值,进行下一次mid后面的比较,直到找到相等的,就可以得到它的位置。
我是一名优秀的程序员,十分优秀!