【问题标题】:Filtering out based on count using Apache Beam使用 Apache Beam 根据计数过滤掉
【发布时间】:2020-10-01 18:10:00
【问题描述】:

我正在使用 Dataflow 和 Apache Beam 处理数据集并将结果存储在包含两列的无头 csv 文件中,如下所示:

A1,a
A2,a
A3,b
A4,a
A5,c
...

我想根据以下两个条件过滤掉某些条目:

1- 在第二列中,如果某个值的出现次数小于N,则删除所有此类行。例如,如果 N=10c 只出现 7 次,那么我希望过滤掉所有这些行。

2- 在第二列中,如果某个值的出现次数超过M,则只保留M许多这样的行并过滤掉其余的行。例如,如果M=1000a 出现了 1200 次,那么我希望过滤掉 200 个此类条目,并将其他 1000 个案例存储在 csv 文件中。

换句话说,我想确保第二列的所有元素出现的次数多于N 且少于M

我的问题是这是否可以通过在 Beam 中使用一些过滤器来实现?还是应该在创建并保存 csv 文件后作为后处理步骤完成?

【问题讨论】:

  • 如果 N=10 且 M=100 而你,例如,“c”出现:a) 99, b) 100 和 c) 101 次。每个案例预计有多少个输出元素?
  • 例如,如果 c 为 99 或 100,则应包括所有情况。当 c 为 101 时,应排除一种情况,而应将其他 100 种情况包含在 csv 文件中(无论它们如何选择)。

标签: google-cloud-dataflow apache-beam dataflow


【解决方案1】:

您可以使用beam.Filter 将与您的范围的下限条件匹配的所有第二列值过滤到 PCollection 中。 然后将该 PCollection(作为 side input)与您的原始 PCollection 相关联,以过滤掉所有需要排除的行。 至于上限,由于您希望保留任何上限数量的元素而不是完全排除它们,因此您应该进行一些后期处理或提出一些组合转换来做到这一点。

使用字数统计的 Python SDK 示例。

class ReadWordsFromText(beam.PTransform):

def __init__(self, file_pattern):
    self._file_pattern = file_pattern

def expand(self, pcoll):
    return (pcoll.pipeline
            | beam.io.ReadFromText(self._file_pattern)
            | beam.FlatMap(lambda line: re.findall(r'[\w\']+', line.strip(), re.UNICODE)))

p = beam.Pipeline()
words = (p 
     | 'read' >> ReadWordsFromText('gs://apache-beam-samples/shakespeare/kinglear.txt')
     | "lower" >> beam.Map(lambda word: word.lower()))
import random
# Assume this is the data PCollection you want to do filter on.
data = words | beam.Map(lambda word: (word, random.randint(1, 101)))
counts = (words 
      | 'count' >> beam.combiners.Count.PerElement())
words_with_counts_bigger_than_100 = counts | beam.Filter(lambda count: count[1] > 100) | beam.Map(lambda count: count[0])

现在你得到一个像这样的pcollection

def cross_join(left, rights):
    for x in rights:
        if left[0] == x:
           yield (left, x)
data_with_word_counts_bigger_than_100 = data | beam.FlatMap(cross_join, rights=beam.pvalue.AsIter(words_with_counts_bigger_than_100))

现在你从数据集中过滤掉了低于下限的元素,得到

注意('king', 66)中的66是我输入的假随机数据。

要使用此类可视化进行调试,您可以使用交互式光束。您可以按照instructions 设置自己的笔记本运行时;或者您可以使用Google Dataflow Notebooks 提供的托管解决方案。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-09-26
    • 1970-01-01
    • 2021-08-19
    相关资源
    最近更新 更多