【发布时间】:2015-01-24 00:29:09
【问题描述】:
我正在尝试使用 Spark Streaming 1.2.0。在某些时候,我按键对流数据进行分组,然后对它们进行一些操作。
以下是一段测试代码:
...
JavaPairDStream<Integer, Iterable<Integer>> grouped = mapped.groupByKey();
JavaPairDStream<Integer, Integer> results = grouped.mapToPair(
new PairFunction<Tuple2<Integer, Iterable<Integer>>, Integer, Integer>() {
@Override
public Tuple2<Integer, Integer> call(Tuple2<Integer, Iterable<Integer>> tp) throws Exception {
TaskContext tc = TaskContext.get();
String ip = InetAddress.getLocalHost().getHostAddress();
int key = tp._1();
System.out.println(ip + ": Partition: " + tc.partitionId() + "\tKey: " + key);
return new Tuple2<>(key, 1);
}
});
results.print();
mapped 是一个 JavaPairDStream,包装了一个虚拟接收器,每秒存储一个整数数组。
我在有两个从属服务器的集群上运行这个应用程序,每个从属服务器有 2 个内核。 当我检查打印输出时,似乎没有将分区永久分配给节点(或以“粘性”方式)。他们经常在两个节点之间移动。这给我带来了麻烦。
在我的实际应用程序中,我需要为每个分区加载相当大量的地理数据。这些地理数据将用于处理流中的数据。我只能为每个分区加载部分地理数据集。如果分区在节点之间移动,我也必须移动地理数据,这可能非常昂贵。
有没有办法让分区保持粘性,即分区 0、1、2、3 留在节点 0,分区 4、5、6、7 留在节点 1?
我尝试将 spark.locality.wait 设置为一个很大的数字,例如 1000000。但它不起作用。
谢谢。
【问题讨论】: