You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink 1.9升级至1.15.2时toRetractStream转换报错及Schema转换咨询

问题解答

1. 错误含义解析

这个报错的核心原因是:Flink 1.15版本中已废弃的toRetractStream不再支持基于LEGACY(旧版)结构化类型的转换逻辑。

  • 从Flink 1.13+开始,Table API的类型系统做了重构,LEGACY模式是为兼容旧版本保留的过渡逻辑,但官方已逐步移除相关支持。
  • 当你调用toRetractStream时,底层仍在尝试用旧的LEGACY类型映射规则转换Table数据,但retract对应的输出逻辑(sink)已经不再兼容这种模式,因此触发该错误。

2. 适配toChangelogStream的实现方案

toChangelogStream是官方指定的替代API,它输出包含变更日志(INSERT/UPDATE/DELETE)的流,以下是适配PaymentSubscriptions样例类的两种方案:

方案一:直接利用样例类类型推导

Flink可以自动识别样例类的结构,只需指定泛型即可,同时通过ChangelogMode过滤出原代码需要的有效数据(对应原逻辑中filter(x => x._1)):

import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.types.RowKind

val paymentSubscriptionsStream: DataStream[PaymentSubscriptions] =
  sTableEnv
    .toChangelogStream[PaymentSubscriptions](
      joinedStreams,
      // 只保留INSERT和UPDATE_AFTER类型的变更,过滤掉DELETE/UPDATE_BEFORE事件
      org.apache.flink.table.api.ChangelogMode.newBuilder()
        .addContainedKind(RowKind.INSERT)
        .addContainedKind(RowKind.UPDATE_AFTER)
        .build()
    )

方案二:显式定义Schema(更可控)

如果需要精确控制字段类型、映射关系,可以手动定义Schema,确保和样例类完全匹配:

import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.flink.table.api.Schema
import org.apache.flink.table.api.DataTypes
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.types.RowKind

// 构建与PaymentSubscriptions字段完全对应的Schema
val schema = Schema.newBuilder()
  .column("company_id", DataTypes.BIGINT())
  .column("company_code", DataTypes.STRING())
  .column("user_access_id", DataTypes.STRING())
  .column("benevity_account_alias", DataTypes.STRING())
  .column("subscription_id", DataTypes.BIGINT())
  .column("subscription_amount", DataTypes.DECIMAL(38, 18)) // 根据实际业务调整精度
  .column("subscription_name", DataTypes.STRING())
  .column("transaction_reference_number", DataTypes.STRING())
  .column("transaction_source", DataTypes.STRING())
  .column("subscription_created_timestamp", DataTypes.STRING())
  .column("date_key", DataTypes.BIGINT())
  .column("exchange_rate_date", DataTypes.BIGINT())
  .column("recipient", DataTypes.STRING())
  .column("portfolio_id", DataTypes.BIGINT())
  .column("pledge_site_pledge_id", DataTypes.STRING())
  .column("match_cap", DataTypes.DECIMAL(38, 18)) // 根据实际业务调整精度
  .column("match_opt_out", DataTypes.STRING())
  .column("pledge_source", DataTypes.STRING())
  .column("pledge_type", DataTypes.STRING())
  .column("pledge_status", DataTypes.STRING())
  .column("transaction_currency", DataTypes.STRING())
  .column("exchange_rate_code", DataTypes.STRING())
  .column("deactivated_flag", DataTypes.INT())
  .column("ts", DataTypes.BIGINT())
  .build()

val paymentSubscriptionsStream: DataStream[PaymentSubscriptions] =
  sTableEnv
    .toChangelogStream[PaymentSubscriptions](
      joinedStreams,
      schema,
      org.apache.flink.table.api.ChangelogMode.newBuilder()
        .addContainedKind(RowKind.INSERT)
        .addContainedKind(RowKind.UPDATE_AFTER)
        .build()
    )

关键注意事项

  • 原代码的filter(x => x._1)是过滤掉撤回事件(DELETE/UPDATE_BEFORE),用ChangelogMode指定保留的事件类型后,无需额外过滤即可得到有效数据。
  • 必须保证样例类字段名与Table中的字段名完全一致(大小写敏感),若不一致可在Schema中通过columnByExpression重命名字段。
  • 对于BigDecimal类型,显式定义Schema时要指定精度和小数位,避免自动推导出现类型不匹配问题。

内容的提问来源于stack exchange,提问作者David Ferguson

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 07:01:02