【发布时间】:2021-08-10 09:25:57
【问题描述】:
我有一个 Dataflow 流式传输作业。我正在使用 BigqueryIO.write 库将行插入 BigQuery 表。 BQ表中有一个列,应该是存储行创建时间戳的。我需要使用 SQL 函数“CURRENT_TIMESTAMP()”来更新该列的值。
我无法使用任何 java 库(如 Instant.now())来获取当前时间戳。因为这将在作业执行期间得出该值。我正在使用 BigQuery 加载作业,其触发频率为 10 分钟。因此,如果我使用任何 java 库来派生时间戳,那么它将不会返回预期的输出。
我在 BigqueryIO.write 中找不到任何方法,该方法将任何 SQL 函数作为输入。那么这个问题有什么解决办法呢?
【问题讨论】:
-
使用行创建时间戳,您是指生成元素的那一刻吗?如果是这样,您可以对 DoFn 中的上下文使用
.timestamp()方法。这应该返回元素本身的时间戳。 -
@Iñigo,timestamp() 方法无济于事。因为它会在 Dataflow 作业执行期间尝试构建时间戳。但正如我在描述中提到的,我需要使用 BigqueryIO File Load 方法将数据插入 BQ 表。触发频率为 10 分钟,这意味着实际 BQ 插入将在 Dataflow 作业执行后 20 分钟(或更多根据数据量可能会分成多个批次)发生。
-
c.timestamp()(c是上下文)在执行期间不会尝试构建时间戳,但它将是元素“创建”的时间戳。例如,如果元素是从 PubSub 读取的消息,c.timestamp()将是该消息的发布时间。无论如何,不确定这是否适用于您的情况。也许使用withFormatFunction并在那里添加时间戳?
标签: google-bigquery bigdata streaming google-cloud-dataflow apache-beam