【问题标题】:How to limit number of lines per file written using FileIO如何限制使用 FileIO 写入的每个文件的行数
【发布时间】:2020-08-28 00:06:37
【问题描述】:

有没有一种方法可以使用 TextIO 或 FileIO 来限制每个写入分片中的行数?

例子:

  1. 从 Big Query 中读取行 - 批处理作业(例如,结果为 19500 行)。
  2. 进行一些转换。
  3. 将文件写入 Google Cloud 存储(19 个文件,每个文件限制为 1000 条记录,一个文件有 500 条记录)。
  4. 触发 Cloud Function 以针对 GCS 中的每个文件向外部 API 发出 POST 请求。

这是我目前尝试做的,但不起作用(尝试限制每个文件 1000 行):

BQ_DATA = p | 'read_bq_view' >> beam.io.Read(
        beam.io.BigQuerySource(query=query,
                               use_standard_sql=True)) | beam.Map(json.dumps)

BQ_DATA | beam.WindowInto(GlobalWindows(), Repeatedly(trigger=AfterCount(1000)),
                              accumulation_mode=AccumulationMode.DISCARDING)
        | WriteToFiles(path='fileio', destination="csv")

我在概念上是错误的还是有其他方法可以实现这一点?

【问题讨论】:

  • 你想任意分配文件中的行数还是基于某种逻辑?另外,你能详细说明一下用例吗?
  • @JayadeepJayaraman 是的,我想任意分配文件中的行数,其中每个文件限制为 1000 行。我已经编辑了我的问题,添加了有关实际用例的更多详细信息。
  • 这是批处理作业吗?批处理作业中会忽略触发器。
  • 是的,它是批量的。感谢彼得让我知道。

标签: google-cloud-dataflow apache-beam


【解决方案1】:

您可以在 ParDo 中实现写入 GCS 步骤并限制要包含在“批处理”中的元素数量,如下所示:

from apache_beam.io import filesystems

class WriteToGcsWithRowLimit(beam.DoFn):
  def __init__(self, row_size=1000):
    self.row_size = row_size
    self.rows = []

  def finish_bundle(self):
     if len(self.rows) > 0:
        self._write_file()

  def process(self, element):
    self.rows.append(element)
    if len(self.rows) >= self.row_size:
        self._write_file()

  def _write_file(self):
    from time import time
    new_file = 'gs://bucket/file-{}.csv'.format(time())
    writer = filesystems.FileSystems.create(path=new_file)
    writer.write(self.rows) # may need to format
    self.rows = []
    writer.close()
BQ_DATA  | beam.ParDo(WriteToGcsWithRowLimit())

请注意,这不会创建任何少于 1000 行的文件,但您可以更改 process 中的逻辑来执行此操作。

(编辑 1 处理余数)

(编辑 2 以停止使用计数器,因为文件将被覆盖)

【讨论】:

  • 成功了,谢谢彼得。我刚刚将writer.write 格式格式化为for row in self.rows: writer.write((str(row) + "\n").encode(encoding='UTF-8'))
  • @HazemSayed 真棒:)
  • 不幸的是,该解决方案仅适用于总共 1000 个元素。这意味着如果我的 Pcollection 有 1001,它将只处理 1000。一个元素被丢弃。这里有任何提示@peter-kim
  • @HazemSayed 我刚刚意识到使用计数器并不是一个理想的解决方案,因为它会覆盖文件。
  • 您可以使用 BatchElements beam.apache.org/releases/pydoc/2.22.0/… 后跟一个写入整个批次的 DoFn。还需要注意从丢弃的包中丢弃文件(参见 Beam 的实现 github.com/apache/beam/blob/master/sdks/python/apache_beam/io/…)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-07-08
  • 1970-01-01
  • 1970-01-01
  • 2016-05-21
  • 1970-01-01
  • 1970-01-01
  • 2018-05-12
相关资源
最近更新 更多