【问题标题】:Apache Beam SQLTransform: Cannot call getSchema when there is no schemaApache Beam SQLTransform:没有架构时无法调用 getSchema
【发布时间】:2019-07-03 22:22:07
【问题描述】:

我正在尝试在PCollection<Object> 上应用 SQLTransform。这里,CustomSource 转换在运行时生成一个 Pojo。因此,在编译时不知道运行 SQLTransform 的 Object 的类型。

        Pipeline p = Pipeline.create(options);

        PCollection<Object> objs = p.apply(new CustomSource());

        Schema type = Schema.builder().addInt32Field("c1").addStringField("c2").addDoubleField("c3").build();
        PCollectionTuple.of(new TupleTag<>("somedata"), objs).apply(SqlTransform.query("SELECT c1 FROM somedata"))
                .setSchema(type, SerializableFunctions.identity(), SerializableFunctions.identity());
        p.run().waitUntilFinish();

我已经使用setSchemaSQLTransform 提供了架构,但我收到了一个错误,即

java.lang.IllegalStateException: Cannot call getSchema when there is no schema
    at org.apache.beam.sdk.values.PCollection.getSchema(PCollection.java:328)
PCollection.java:328
    at org.apache.beam.sdk.extensions.sql.impl.schema.BeamPCollectionTable.<init>(BeamPCollectionTable.java:34)

是否可以在运行时生成 Pojo 对象并通过向转换提供架构信息来对其运行 sqltransforms ?

这是 CustomSource 类供参考:

import java.util.HashMap;
import java.util.Map;

import com.beaconinside.messages.PojoGenerator;

import org.apache.beam.sdk.coders.StringUtf8Coder;
import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.values.PBegin;
import org.apache.beam.sdk.values.PCollection;

import javassist.CannotCompileException;
import javassist.NotFoundException;

public class CustomSource extends PTransform<PBegin, PCollection<Object>> {

    Map<String, Class<?>> props;
    Class<?> clazz;
    String data = "{\"c1\": 1, \"c2\": \"row\", \"c3\": 2.0}";

    public CustomSource() throws NotFoundException, CannotCompileException {
        props = new HashMap<String, Class<?>>();
        props.put("c1", Integer.class);
        props.put("c2", String.class);
        props.put("c3", Double.class);
        clazz = PojoGenerator.generate("net.javaforge.blog.javassist.PojoGenerated", props);
    }

    @Override
    public PCollection<Object> expand(PBegin input) {
        return input.apply(Create.of(data)).setCoder(StringUtf8Coder.of()).apply(new SensorSource(clazz, props));
        // return input.apply(Create.of(data));
    }

}

【问题讨论】:

  • 请指点一下?

标签: apache-beam beam-sql


【解决方案1】:

我认为您的 setSchema 只是来自SQLTransform 的输出 PCollection 的设置模式。您还应该在 PCollection&lt;Object&gt; objs 上设置架构。

【讨论】:

    【解决方案2】:

    上面的答案是对的,PCollection&lt;Object&gt; 也应该调用setSchema 来定义输入数据模式和行对象转换函数。如果您有多个 PCollection 来构建 PCollectionTuple,则 PCollection 应分别调用 setSchema。 PCollectionTuple 不需要调用 setSchema,因为可以从 SQL 命令中推断出输出模式。

    【讨论】:

      【解决方案3】:

      如下使用 setRowSchema

      PCollection<Row> testApps = PBegin.in(p).apply(Create.of(row1,row2,row3).withCoder(RowCoder.of(appSchema)))
                      .setRowSchema(appSchema);
      

      【讨论】:

        猜你喜欢
        • 2023-01-08
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2012-11-22
        • 1970-01-01
        • 1970-01-01
        • 2018-08-11
        • 1970-01-01
        相关资源
        最近更新 更多