【问题标题】:How to write Custom coder for TableRow wrapper with addition parameter如何使用附加参数为 TableRow 包装器编写自定义编码器
【发布时间】:2022-11-22 05:41:34
【问题描述】:

我正在从事 Apachebeam 管道项目,以从 GCS 存储桶中的 .TSV 文件读取数据,将数据转换为 BigQuery 行并将其写入 GCS 中的表。

我还必须根据输入文件中的值来确定表名。因此我创建了一个包装类如下

public class TableAndRow implements Serializable {

    @Nullable
    public String tab_name;
    @Nullable
    public TableRow row;
}

我在 DoFn 方法中将文件数据转换为包装类对象,但由于编码器问题,系统抛出错误。请帮助如何为此类包装类编写编码器。

我的管道代码如下所示

    lines.apply("Convert Each line to TableRow", ParDo.of(new DoFn<String, TableAndRow>() {

            @ProcessElement 
            public void processElement(ProcessContext c) {
                TableAndRow output_row = new TableAndRow();
                //TableRow output_row = new TableRow();
                String[] arr = c.element().split("\t");
                output_row = getRow(arr);
                
                c.output(output_row);
            }

我尝试使用 AcroCoder 但导致以下错误 org.apache.avro.UnresolvedUnionException: Not in union ["null",{"type":"record","name":"TableRow","namespace":"com. google.api.services.bigquery.model","fields":[{"name":"f","type":{"type":"array","items":{"type":"record" “名称”:“TableCell”,“字段”:[{“名称”:“v”,“类型”:{“类型”:“记录”,“名称”:“对象”,“命名空间”:“java .lang","fields":[]}},{"name":"jsonFactory","type":{"type":"record","name":"JsonFactory","namespace":"com. google.api.client.json","fields":[]}},{"name":"unknownFields","type":{"type":"map","values":"java.lang.Object "}},{"name":"classInfo","type":{"type":"record","name":"ClassInfo","namespace":"com.google.api.client.util", “字段”:[{“名称”:“clazz”,“类型”:{“类型”:“记录”,“名称”:“类”,“命名空间”:“java.lang”,“字段”:[ ]}},{"name":"ignoreCase","type":"boolean"},{"name":"nameToFieldInfoMap","type":{"type":"map","values":{" type":"record","name":"FieldInfo","fields":[{"name":"isPrimitive","type":"boolean"},{"name":"field ","type":{"type":"record","name":"Field","namespace":"java.lang.reflect","fields":[]}},{"name":" setters","type":{"type":"array","items":{"type":"record","name":"Method","namespace":"java.lang.reflect"," fields":[]},"java-class":"[Ljava.lang.reflect.Method;"}},{"name":"name","type":"string"}]}}},{ "name":"names","type":{"type":"array","items":"string","java-class":"java.util.List"}}]}}]}, "java-class":"java.util.List"}},{"name":"jsonFactory","type":"com.google.api.client.json.JsonFactory"},{"name":" unknownFields","type":{"type":"map","values":"java.lang.Object"}},{"name":"classInfo","type":"com.google.api. client.util.ClassInfo"}]}]: GenericData{classInfo=[f], {eventType=detail-page-view, visitorId=89395430694564180440746546053353344574, eventTime=2022-10-11 23:40:33, experimentIds=bloomreach, productDetails .product.id=;BSH15730;;;;125=我的商店:n^|附近商店:n^|DC商店:n|139=::hash::0|157=::hash::0|165= ::hash::0|169=::hash::0|170=::hash::0|282=内部搜索|283=::hash::0|284=::hash::0| 285=1:1|286=是|287=否|288=2+ 天|289=明天|291=o2 传感器 1|293=否|294=否|295=::hash::0|296=: :hash::0|297=::hash::0, userInfo.userId=1, userInfo.ipAddress=97003, userInfo.userAgent=Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/105.0.0.0 Safari/537.36, userInfo.directUserRequest=1, uri=https://www.napaonline.com/en/c/brakes, referrerUri=1}}

【问题讨论】:

  • @Deniz Saner 我看到了您的一篇帖子,您在其中为类似的包装类编写了编码器。你能帮忙吗

标签: pipeline apache-beam tablerow encoder beam


【解决方案1】:

我认为您可以使用 MapElementsString 行直接映射到 TableRow,在这种情况下将推断出编码器,例如:

    @Test
    public void testMapToTableRow() {
        final List<String> input = Arrays.asList("toto,tata");

        pipeline.apply("Read inputs", Create.of(input))
                .apply("Convert Each line to TableRow", MapElements.into(of(TableRow.class)).via(this::toRow))
                .apply("Print element", MapElements.into(of(TableRow.class)).via(this::printRow));

        pipeline.run().waitUntilFinish();
    }

    private TableRow toRow(final String element) {
        final TableRow row = new TableRow();
        
        // Put your logic here with the split on input and the instantiation of your TableRow object.
        // row.set("name", "FFF");
        return row;
    }

    private TableRow printRow(final TableRow row) {
        System.out.println(row);

        return row;
    }

我为我的示例使用了单元测试。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2013-05-12
    • 2021-04-19
    • 2023-03-04
    • 1970-01-01
    • 1970-01-01
    • 2011-07-25
    • 2013-04-08
    • 2012-03-15
    相关资源
    最近更新 更多