我已经借助 Kafka 按键将数据排序到我的 Spark Streaming 分区中,即在一个节点上找到的键在任何其他节点上都找不到。
\n\n我想使用 redis 及其incrby(增量)命令作为状态引擎,并减少发送到 redis 的请求数量,我想通过对每个工作节点本身进行字数统计来部分减少我的数据。(关键是标签+时间戳来从字数统计中获取我的功能)。\n我想避免洗牌并让 Redis 负责跨工作节点添加数据。
即使我检查数据在工作节点之间干净地分割,.reduce(_ + _)(Scala 语法)也需要很长时间(映射任务需要几秒而不是亚秒),因为 HashPartitioner 似乎将我的数据洗牌到随机节点以添加它那里。
如何使用 Spark Streaming 在每个分区器上编写一个简单的字数减少而不触发 Scala 中的洗牌步骤?
\n\n注意 DStream 对象缺少一些 RDD 方法,这些方法只能通过该transform方法获得。
看来我或许可以使用combineByKey。我想跳过这mergeCombiners()一步,而是将累积的元组保留在原处。\n《Learning Spark》一书神秘地说:
\n\n\n如果我们知道我们的数据不会从中受益,我们可以在combineByKey()中禁用地图端聚合。例如,groupByKey() 禁用映射端聚合,因为聚合函数(附加到列表)不会节省任何空间。如果我们想禁用映射端联合,我们需要指定分区器;现在,您可以通过传递 rdd.partitioner 在源 RDD 上使用分区器。
\n
https://www.safaribooksonline.com/library/view/learning-spark/9781449359034/ch04.html
\n\n然后,这本书仍然没有提供如何执行此操作的语法,到目前为止我也没有在谷歌上找到任何运气。
\n\n更糟糕的是,据我所知,Spark Streaming 中没有为 DStream RDD 设置分区器,所以我不知道如何为ombineByKey 提供分区器,而不会最终打乱数据。
\n\n另外,“地图端”实际上意味着什么以及mapSideCombine = false到底会产生什么后果?
的 scala 实现combineByKey位于\n …