【发布时间】:2020-06-10 13:28:51
【问题描述】:
嗨,
我正在使用 Beam 从 BQ 表中读取数据,发现使用 SerializableFunction 的 read() 比 readTableRows() 具有更好的性能。按照https://beam.apache.org/releases/javadoc/2.20.0/org/apache/beam/sdk/io/gcp/bigquery/BigQueryIO.html#read-org.apache.beam.sdk.transforms.SerializableFunction-的示例
我的 Big Query 列是:
|Field name | Field type|
|Date_Time | TIMESTAMP |
|Simple_Id | STRING |
|A_Price | NUMERIC |
我的代码如下:
公共类 ConvertBQSchemaRecordToProtoDataFn 实现 SerializableFunction {
@Override
public ProtoValueType apply(SchemaAndRecord schemaAndRecord) {
GenericRecord avroRecord = schemaAndRecord.getRecord();
long dateTimeMillis = (Long) avroRecord.get("Date_Time");
String simpleId = avroRecord.get("Simple_Id").toString();
double aPrice = convertToDouble(avroRecord.get("A_Price").toString());
long 和 String 都可以。但是,当我尝试转换 NUMERIC 类型时,GenericRecord(来自调试器)将其显示为您无法转换的 HeapByteBuffer。我不确定如何获取“A_Price”的值:
调用管道代码如下:
PCollection<ProtoValueType> protoData =
pipeline.apply("BigQuery Read",
BigQueryIO.read(new ConvertBQSchemaRecordToProtoDataFn())
.fromQuery(sqlQuery)
.usingStandardSql()
.withCoder(ProtoCoder.of(ProtoValueType.class)));
我不确定是否使用了 Coder。 ProtoValueType 是一个 protobuf 生成的绑定类。
我的问题是:如何从 GenericRecord(我认为是 Avro 对象)中获取 NUMERIC 类型的值?
任何帮助表示赞赏。我可以使用 readTableRows() 获取该行,它都以字符串形式返回,所以我不想理解该方法。
【问题讨论】:
-
只是为了澄清,您正在从 BigQuery 中读取,您在哪里写入输出?
-
您好 Alexandre,我正在从 BigQuery 读取数据并将行转换为 protobuf 对象,我将传递给另一个函数(例如,平均 aPrice 值)。输出可以是平均值,也可以是其他东西(仍在编写管道)。
标签: java google-bigquery apache-beam