【问题标题】:Ranking pcollection elements对 pcollection 元素进行排名
【发布时间】:2018-01-16 15:20:50
【问题描述】:

我正在使用 Google DataFlow Java SDK 2.2.0。用例如下:

PCollection pEmployees:员工及对应部门名称。最多可包含 1000 万个元素。

PCollection pDepartments:部门名称和每个部门要发布的元素数量。将包含数百个元素。

任务:根据 pDepartments 中所有部门的部门编号从 pEmployees 中收集元素。这将是一个大集合(最多几十万个元素或几 GB)。

我们不能在此处使用 Top 转换,因为它会在 pEmployee 上一次工作一个,而我们有多个部门,而且在 PCollection 中也有。我们可以为 pEmployees 中的每个元素分配一个行号,将其与 pDepartments 连接,并从 pDepartments 中过滤 row_number > target number 的记录。这将需要一个全球排名。

问题:我们如何将排名/行号分配给 pcollection 中的元素?

【问题讨论】:

  • 我是否理解正确,您想从每个部门中选择不同数量的员工?在一个部门内,选择应该是任意的,还是例如“部门内薪酬最高的前 N ​​名员工”?
  • 是的,每个部门需要招聘的员工人数不同。目前,该用例不需要在部门内进行有序选择,但如果有这样就很好了。

标签: google-cloud-dataflow apache-beam


【解决方案1】:

这非常接近Sample 转换,但不完全是,因为当用作.perKey() 时,它对所有键应用相同的阈值。一般来说,Beam 目前不支持 per-key 组合不同的组合函数参数。

我建议通过使用CoGroupByKey 来模拟它,加入pEmployeespDepartments 并获取包含部门名称、N = 元素数和该部门所有员工的元组 (CoGbkResult)。然后简单地遍历员工并发出第一个 N 并丢弃其余的。

【讨论】:

  • 这听起来是个不错的方法,但是在CoGroupByKey 之后,由于每个部门的员工数量,一些元组的大小可能会很大。如果我们在ParDo 中处理它,可能会导致管道倾斜。这会被推荐用于生产吗?
  • 我认为有几百万个元素分布(即使不均匀)超过几百个键,没有什么可担心的。我建议尝试一下,看看性能是否不够,这是瓶颈。如果是这样,我可以推荐一些替代但更复杂的方法。
  • 请注意,具有公共键的连接集合的元素不需要放入内存中,即使在单个键中也是如此。
  • 谢谢!是否有与此功能相关的任何文档?
  • 不确定您的意思? CoGroupByKey 有 javadoc。
猜你喜欢
  • 2011-06-22
  • 1970-01-01
  • 1970-01-01
  • 2015-08-23
  • 2020-11-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多