【问题标题】:Scala Option Types not recognized in apache flink table apiapache flink table api中无法识别的Scala选项类型
【发布时间】:2022-10-12 22:25:06
【问题描述】:

我正在构建一个从 kafka 读取数据的 flink 应用程序 主题,应用一些转换并写入 Iceberg 表。

我从 kafka 主题(在 json 中)读取数据并使用 circe 解码 到 scala 案例类,其中包含 scala 选项值。 数据流上的所有转换都可以正常工作。

案例类如下所示

Event(app_name: Option[String], service_name: Option[String], ......)

但是当我尝试将流转换为表以写入冰山表时 由于案例类,列被转换为原始类型,如图所示 以下。

table.printSchema()

service_name RAW('scala.Option', '...'),
conversion_id RAW('scala.Option', '...'),
......

并且表写入失败如下。

Query schema: [app_name: RAW('scala.Option', '...'), .........
Sink schema: [app_name: STRING, .......

flink table api 是否支持带有选项值的 scala 案例类? https://nightlies.apache.org/flink/flink-docs-master/docs/dev/datastream/fault-tolerance/serialization/types_serialization/#special-types

我发现本文档中的数据流支持它。

有没有办法在 Table API 中做到这一点。

在此先感谢您的帮助..

【问题讨论】:

    标签: scala apache-flink flink-streaming flink-sql iceberg


    【解决方案1】:

    Table API 的类型系统比 DataStream API 的类型系统更严格。不受支持的类会立即被视为黑盒类型RAW。这允许对象仍然通过 API,但可能不是每个连接器都支持它。

    从例外情况来看,您似乎使用app_name: STRING 声明了接收器表,所以我想您可以使用字符串表示形式。如果是这种情况,我建议实现一个user-defined function 来执行到字符串的转换。

    【讨论】:

      【解决方案2】:

      我遇到了完全相同的 Option 被解析为 RAW 的问题,并找到了另一种您可能感兴趣的解决方法:

      TL;博士:
      而不是使用返回Option.get 并努力将返回类型声明为Types.OPTION(Types.INSTANT)这是行不通的,我改为使用.getOrElse("key", null) 并将类型声明为常规类型。 Table API 然后识别列类型,创建一个可为空的列并正确解释空值。然后我可以使用IS NOT NULL 过滤这些行。

      详细示例:

      对我来说,它从一个自定义地图功能开始,我在其中解压缩可能缺少某些字段的数据:

      class CustomExtractor extends RichMapFunction[MyInitialRecordClass, Row] {
        def map(in: MyInitialRecordClass): Row = {
          Row.ofKind(
            RowKind.INSERT,
            in._id
            in.name
            in.getOrElse("time", null) // This here did the trick instead of using .get and returning an option
          )
        }
      }
      

      然后我像这样显式地声明一个返回类型。

      val stream: DataStream[Row] = data
        .map(new CustomExtractor())
        .returns(
          Types.ROW_NAMED(
            Array("id", "name", "time"),
            Types.STRING,
            Types.STRING,
            Types.INSTANT
          )
      
      val table = tableEnv.fromDataStream(stream)
      table.printSchema()
      // (
      //   `id` STRING,
      //   `name` STRING,
      //   `time` TIMESTAMP_LTZ(9),
      // )
      
      tableEnv.createTemporaryView("MyTable", table)
      
      tableEnv
        .executeSql("""
            |SELECT
            | id, name, time
            |FROM MyTable""".stripMargin)
        .print()
      // +----+------------------+-------------------+---------------------+
      // | op |               id |              name |                time |
      // +----+------------------+-------------------+---------------------+
      // | +I |              foo |               bar |              <Null> |
      // | +I |             spam |               ham | 2022-10-12 11:32:06 |
      

      这至少对我来说正是我想要的。我对 Flink 非常陌生,如果这里的专业人士认为这是可行的或非常骇人听闻,我会很好奇。

      使用 scala 2.12.15 和 flink 1.15.2

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2013-07-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多