【问题标题】:No support for the type of the given DataStream: GenericType<org.apache.flink.types.Row> Flink Cassandra不支持给定 DataStream 的类型:GenericType<org.apache.flink.types.Row> Flink Cassandra
【发布时间】:2021-09-17 20:13:45
【问题描述】:

我想向 Cassandra 写入行流。 首先,我将 Avro 流转换为行流。编译时没有显示错误。 请参阅下面的代码:(KafkaConsumer 和 CassandraSink 在其他工作中分别工作正常)

StreamExecutionEnvironment environment =  StreamExecutionEnvironment.getExecutionEnvironment();

// Initialize KafkaConsumer
FlinkKafkaConsumer010 kafkaConsumer = KafkaConnection.getKafkaConsumer(AvroSchemaClass.class, inTopic, schemaRegistryUrl, properties);

// Set KafkaConsumer as source
DataStream<AvroSchemaClass> avroInputStream = environment.addSource(kafkaConsumer);

// converting avro message to flink's row datatype.
// see https://ci.apache.org/projects/flink/flink-docs-master/api/java/org/apache/flink/formats/avro/AvroRowDeserializationSchema.html
AvroRowDeserializationSchema avroToRow = new AvroRowDeserializationSchema(AvroSchemaClass.class);
DataStream<Row> rowInputStream = avroInputStream.map(new MapFunction<Orders_value, Row>() {
                @Override
                public Row map(AvroSchemaClass orders_value) throws Exception {
                    return avroToRow.deserialize(orders_value.toByteBuffer().array());
                }
            });

// Example transformation
DataStream<Row> rowOutputStream = rowInputStream.filter(row -> country.equals(row.getField(7).toString()));
       
CassandraSink streamSink = CassandraConnection.getSink(rowOutputStream,
                    cassandraURL,
                    cassandraPort,
                    cassandraCluster,
                    cassandraUser,
                    cassandraPass,
                    insertQuery);
streamSink.name("Write something to Cassandra");

environment.execute();

但是当我在flink中运行作业时,出现如下错误:

java.lang.IllegalArgumentException: No support for the type of the given DataStream: GenericType<org.apache.flink.types.Row>
        at org.apache.flink.streaming.connectors.cassandra.CassandraSink.addSink(CassandraSink.java:255)
        at servingLayer.CassandraConnection.getSink(CassandraConnection.java:24)
        at speedLayer.KafkaToCassandra.main(KafkaToCassandra.java:84)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
        at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
        at java.base/java.lang.reflect.Method.invoke(Unknown Source)
        at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:355)
        at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:222)
        at org.apache.flink.client.ClientUtils.executeProgram(ClientUtils.java:114)
        at org.apache.flink.client.cli.CliFrontend.executeProgram(CliFrontend.java:812)
        at org.apache.flink.client.cli.CliFrontend.run(CliFrontend.java:246)
        at org.apache.flink.client.cli.CliFrontend.parseAndRun(CliFrontend.java:1054)
        at org.apache.flink.client.cli.CliFrontend.lambda$main$10(CliFrontend.java:1132)
        at org.apache.flink.runtime.security.contexts.NoOpSecurityContext.runSecured(NoOpSecurityContext.java:28)
        at org.apache.flink.client.cli.CliFrontend.main(CliFrontend.java:1132)
java.lang.NullPointerException

数据流类型的特定更改会是解决方案吗?如果是,如何实施? 如果您需要更多信息,请告诉我。

【问题讨论】:

  • 将 Avro 类型转换为 Row 类型的想法是什么?我认为没有这一步你会更好。
  • @Arvid 我在其他工作中使用 Row 类型,所以我希望我可以为仅使用 Row 类型的转换编写通用代码。是否有更好的方法将 Avro 转换为 Row,以便与 CassandraSink 一起使用?

标签: java cassandra apache-flink


【解决方案1】:

似乎CassandraSink 应该支持Row 开箱即用。问题是rowOutputStreamRowTypeInfo 不知何故丢失了,它使用了备用GenericType(Kryo 序列化效率低下)。

AvroRowDeserializationSchema 正在正确返回类型信息,但 DataStream API 没有自动获取。

因此,如果一切都成立,那么解决方法是显式设置 rowIn/OutputStream 的返回类型,如下所示

DataStream<Row> rowInputStream = avroInputStream.map(new MapFunction<Orders_value, Row>() {
            @Override
            public Row map(AvroSchemaClass orders_value) throws Exception {
                return avroToRow.deserialize(orders_value.toByteBuffer().array());
            }
        }).returns(avroToRow.getProducedType());
...
DataStream<Row> rowOutputStream = rowInputStream.filter(row -> country.equals(row.getField(7).toString()))
  .returns(avroToRow.getProducedType())

一般来说,如果您坚持使用一种 API,会更容易。在这种情况下,我建议完全使用 Table API。

【讨论】:

  • 不幸的是,rowInputStream 已经有错误的GenericType 类型。所以问题的根源可能不是filter,而是AvroRowDeserializationSchema或avroSchemaClassOrders_value。我将 avroSchemaClass 更改为 avroSchemaString,但错误仍然存​​在。
  • 感谢调查,我更新了我的答案,希望它现在可以工作。
  • 抱歉回复晚了。谢谢你的更新。 DataStream.returns() 不存在,但我找到了设置返回类型的解决方法:Transformation rowInputStreamTransformation = rowInputStream.getTransformation(); rowInputStreamTransformation.setOutputType(avroToRow.getProducedType()); 不幸的是,我无法确认,它正在工作,因为 Cassandra 没有获取任何值,但 flink 日志中没有任何错误。你对这个问题有什么建议吗?
  • 我再次更新解决方案。
  • 感谢您的更新。长话短说:我可能需要问另一个问题,如何将 ConfluentAvroSchema 转换为 Row:我得到以下 java.io.IOException: Failed to deserialize Avro record. 仅供参考:我使用 Confluent Schema Registry。所以我可能需要使用特定的ConfluentAvroRowDeserializationSchema。我找到了这个 repo github.com/ztore/flink-registry-avro-row-schema。但它的功能需要flink.formats.avro.SchemaCoder.SchemaCoderProvider 作为参数。由于我是一个java初学者,我不知道如何根据函数的需要来实现这个接口。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-09-12
  • 1970-01-01
  • 1970-01-01
  • 2018-09-13
  • 2015-03-02
  • 1970-01-01
  • 2020-12-23
相关资源
最近更新 更多