【发布时间】:2020-08-28 00:06:37
【问题描述】:
有没有一种方法可以使用 TextIO 或 FileIO 来限制每个写入分片中的行数?
例子:
- 从 Big Query 中读取行 - 批处理作业(例如,结果为 19500 行)。
- 进行一些转换。
- 将文件写入 Google Cloud 存储(19 个文件,每个文件限制为 1000 条记录,一个文件有 500 条记录)。
- 触发 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