【问题标题】:How to create PCollection<Row> from PCollection<String> for performing beam SQL Trasforms如何从 PCollection<String> 创建 PCollection<Row> 以执行 beam SQL Trasforms
【发布时间】:2019-07-05 08:10:23
【问题描述】:

我正在尝试实现一个数据管道,它连接来自 Kafka 主题的多个无限源。我能够连接到主题并以PCollection&lt;String&gt; 获取数据,我需要将其转换为PCollection&lt;Row&gt;。我将逗号分隔的字符串拆分为一个数组,并使用模式将其转换为行。但是,如何实现/构建架构并将值动态绑定到它?

即使我为架构构建创建了一个单独的类,有没有办法将字符串数组直接绑定到架构?

下面是我当前的工作代码,它是静态的,每次构建管道时都需要重写,并且它也会根据字段的数量而延长。

final Schema sch1 =
                Schema.builder().addStringField("name").addInt32Field("age").build();

PCollection<KafkaRecord<Long, String>> kafkaDataIn1 = pipeline
  .apply(
    KafkaIO.<Long, String>read()
      .withBootstrapServers("localhost:9092")
      .withTopic("testin1")
      .withKeyDeserializer(LongDeserializer.class)
      .withValueDeserializer(StringDeserializer.class)
      .updateConsumerProperties(
         ImmutableMap.of("group.id", (Object)"test1")));

PCollection<Row> Input1 = kafkaDataIn1.apply(
  ParDo.of(new DoFn<KafkaRecord<Long, String>, Row>() {
    @ProcessElement
    public void processElement(
        ProcessContext processContext,
        final OutputReceiver<Row> emitter) {

          KafkaRecord<Long, String> record = processContext.element();
          final String input = record.getKV().getValue();

          final String[] parts = input.split(",");

          emitter.output(
            Row.withSchema(sch1)
               .addValues(
                   parts[0],
                   Integer.parseInt(parts[1])).build());
        }}))
  .apply("window",
     Window.<Row>into(FixedWindows.of(Duration.standardSeconds(50)))
       .triggering(AfterWatermark.pastEndOfWindow())
       .withAllowedLateness(Duration.ZERO)
       .accumulatingFiredPanes());

Input1.setRowSchema(sch1);

我的期望是以动态/可重用的方式实现与上述代码相同的事情。

【问题讨论】:

    标签: java join apache-kafka apache-beam


    【解决方案1】:

    架构是在 pcollection 上设置的,因此它不是动态的,如果您想懒惰地构建它,那么您需要使用支持它的格式/编码器。 Java 序列化或 json 就是示例。

    也就是说,要从 sql 功能中受益,您还可以使用带有查询字段和其他字段的静态模式,这样静态部分可以让您执行 sql 并且您不会丢失附加数据。

    罗曼

    【讨论】:

    • 感谢您的回复。如果你能提供一个例子,对我理解会更有用。
    • github.com/Talend/component-runtime/blob/master/… 是一个通用编码器,它只是通过 jsonobject(使用 jsonp,但你也可以使用 jackson)。这个想法是为非动态部分使用 + 静态编码器 - 可以建模作为梁行。将两者结合起来,您将获得一个行模式(例如,通用部分只是嵌入在一列中)。对于 json 字符串列非常好,然后使用梁的构建器很容易构建模式 bit.ly/2Y5vFpH
    猜你喜欢
    • 2022-12-31
    • 2023-02-03
    • 2023-04-10
    • 2019-02-09
    • 2021-12-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多