小编use*_*464的帖子

设置自定义编码器和处理参数化类型

我有两个与我在Dataflow管道中遇到的编码器问题有关的问题.

  • 如何为自定义数据类型设置编码器?该类只包含三个项目 - 两个双精度数和另一个参数化属性.我尝试用SerializableCoder注释类型,但我仍然得到错误"com.google.cloud.dataflow.sdk.coders.CannotProvideCoderException:无法根据类接口java.util.Set提供编码器:没有注册CoderFactory为了上课." Set实际上包含参数化的自定义数据类型 - 所以我假设自定义数据类型是问题.我找不到足够的文档/示例正确的方法来做到这一点.请将我指向正确的地方.
  • 即使没有自定义数据类型,每当我尝试切换到参数化版本的Transform函数时,都会导致编码器错误.具体来说,在参数化的复杂变换中,ParDo使用参数化类型,但是当我在ParDo之后对结果PCollection应用Combine.PerKey时,会导致CoderNotFoundException.

关于这两个项目的任何帮助都会有所帮助,因为我现在有点困惑.

google-cloud-platform google-cloud-dataflow

5
推荐指数
1
解决办法
1162
查看次数

用于处理多个Pubsub主题的数据流管道设计

我有一个从Pubsub主题读取的管道(按分钟窗口)并将处理结果写入BigQuery.我想让表格按时间分片,以及数据本身的一些键.BigQueryIO确实通过窗口时间戳为shard提供了选项,但我认为它不提供任何选项来通过输入集合本身的某些键对表进行分片.如果我错过了一些替代方案,请告诉我.

为了克服这个问题,(选项1)我选择使用相同的密钥对源Pubsub主题本身进行分片,因此,设置管道以从多个源读取并按照单独的分支处理它们并将每个分支结果写入由窗口分区的BigQuery时间戳似乎有效.我想知道的是,由于Dataflow中的中间处理步骤在我的情况下可以与源或接收器无关(选项2)如果我继续使用它会使管道更有效(在资源和时间方面)单个Pubsub主题并在BigQuery编写步骤之前添加额外的转换以对集合进行分区,然后写入BigQuery.

选项 - 1 +在读取/写入期间在Pubsub上进行较小的加载,因为即使组合的消息可能适合几百KB - 读取步骤和中间处理在单独的管道中完成(对于Dataflow可能效率不高)

选项 - 2 +管道更清洁 - 分区的附加步骤也读取与我们分区数量相同的集合次数 - 但是收集项目和分区本身的数量非常小 - 所以,这不应该是一个更大的问题

我认为选择2在阅读管道设计原则时更有意义,但我仍然想澄清我正在做的是对的.

google-cloud-platform google-cloud-dataflow

4
推荐指数
1
解决办法
1126
查看次数

使用Avrocoder进行自定义类型和泛型

我正在尝试使用AvroCoder序列化自定义类型,该类型在我的管道中的PCollections中传递.自定义类型有一个通用字段(当前是一个字符串)当我运行管道时,我得到如下的AvroTypeException可能是由于泛型字段.为这个对象构建和传递AvroSchema是解决这个问题的唯一方法吗?

Exception in thread "main" org.apache.avro.AvroTypeException: Unknown type: T
 at org.apache.avro.specific.SpecificData.createSchema(SpecificData.java:255)
 at org.apache.avro.reflect.ReflectData.createSchema(ReflectData.java:514)
 at org.apache.avro.reflect.ReflectData.createFieldSchema(ReflectData.java:593)
 at org.apache.avro.reflect.ReflectData.createSchema(ReflectData.java:472)
 at org.apache.avro.specific.SpecificData.getSchema(SpecificData.java:189)
 at com.google.cloud.dataflow.sdk.coders.AvroCoder.of(AvroCoder.java:116)
Run Code Online (Sandbox Code Playgroud)

我还附上了我的注册表代码以供参考.

pipelineCoderRegistry.registerCoder(GenericTypeClass.class, new CoderFactory() {
    @Override
    public Coder<?> create(List<? extends Coder<?>> componentCoders) {
        return AvroCoder.of(GenericTypeClass.class);
    }

    @Override
    public List<Object> getInstanceComponents(Object value) {
        return Collections.singletonList(((GenericTypeClass<Object>) value).key);
    }
});
Run Code Online (Sandbox Code Playgroud)

avro google-cloud-platform google-cloud-dataflow

4
推荐指数
1
解决办法
1299
查看次数