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
相关产品推荐
相关产品推荐

