如何实现通用Flink-Kafka到Oracle作业的动态JSON反序列化?
通用Flink作业:Kafka动态JSON反序列化适配Oracle表方案
核心思路
因为Kafka Topic的JSON结构在作业启动前未知,而Oracle表结构可以在作业启动后动态获取,所以核心是先拉取Oracle表元数据,再基于元数据生成对应的JSON反序列化规则,同时替换废弃的JSONDeserializationSchema。
1. 替代废弃的JSONDeserializationSchema
Flink 1.14+已废弃JSONDeserializationSchema,官方推荐使用JsonRowDataDeserializationSchema(基于Flink的RowData模型),或自定义DeserializationSchema结合Jackson实现灵活解析。
2. 动态获取Oracle表结构
作业启动后,通过JDBC连接Oracle,查询系统表获取目标表的字段名、数据类型,再映射为Flink的RowType(这是后续反序列化的核心依据)。
示例代码片段:
// 从Oracle获取表结构,转换为Flink RowType public static RowType getOracleTableRowType(ParameterTool params) throws SQLException { String jdbcUrl = params.get("oracle.jdbc.url"); String username = params.get("oracle.username"); String password = params.get("oracle.password"); String tableName = params.get("oracle.table"); List<RowField> fields = new ArrayList<>(); try (Connection conn = DriverManager.getConnection(jdbcUrl, username, password); Statement stmt = conn.createStatement()) { // 查询Oracle表字段元数据 ResultSet rs = stmt.executeQuery( "SELECT COLUMN_NAME, DATA_TYPE, DATA_LENGTH, DATA_PRECISION, DATA_SCALE " + "FROM USER_TAB_COLUMNS WHERE TABLE_NAME = '" + tableName.toUpperCase() + "'" ); while (rs.next()) { String colName = rs.getString("COLUMN_NAME").toLowerCase(); String dataType = rs.getString("DATA_TYPE"); // 映射Oracle类型到Flink类型 DataType flinkType = mapOracleTypeToFlink(dataType, rs.getInt("DATA_PRECISION"), rs.getInt("DATA_SCALE")); fields.add(new RowField(colName, flinkType)); } } return RowType.of(fields.stream().map(RowField::getType).toArray(DataType[]::new), fields.stream().map(RowField::getName).toArray(String[]::new)); } // 类型映射示例,可根据实际需求扩展 private static DataType mapOracleTypeToFlink(String oracleType, int precision, int scale) { return switch (oracleType.toUpperCase()) { case "VARCHAR2", "CHAR" -> DataTypes.STRING(); case "NUMBER" -> { if (scale == 0) { if (precision <= 10) yield DataTypes.INT(); else yield DataTypes.BIGINT(); } else { yield DataTypes.DECIMAL(precision, scale); } } case "DATE", "TIMESTAMP" -> DataTypes.TIMESTAMP(3); default -> DataTypes.STRING(); // 默认转字符串兼容 }; }
3. 动态生成JSON反序列化器
基于上面获取的RowType,可以直接使用JsonRowDataDeserializationSchema实现动态反序列化,或者自定义逻辑处理更复杂的场景。
方式一:使用JsonRowDataDeserializationSchema(推荐)
ParameterTool params = ParameterTool.fromArgs(args); RowType targetRowType = getOracleTableRowType(params); // 构建动态反序列化器 JsonRowDataDeserializationSchema deserializationSchema = new JsonRowDataDeserializationSchema.Builder(targetRowType) .ignoreParseErrors(true) // 忽略解析错误,避免作业崩溃 .failOnMissingField(false) // 允许JSON缺失表中字段,后续可补默认值 .build(); // 创建Kafka Source KafkaSource<RowData> kafkaSource = KafkaSource.<RowData>builder() .setBootstrapServers(params.get("kafka.bootstrap.servers")) .setTopics(params.get("kafka.topic")) .setGroupId("flink-oracle-sync-group") .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(deserializationSchema) .build();
方式二:自定义DeserializationSchema(Schema-less场景)
如果需要更灵活的处理(比如JSON结构和Oracle表不完全匹配),可以自定义反序列化器,先将JSON转为Map<String, Object>,再根据Oracle表结构组装成Row:
public class DynamicJsonDeserializationSchema implements DeserializationSchema<Row> { private final RowType targetRowType; private transient ObjectMapper objectMapper; public DynamicJsonDeserializationSchema(RowType targetRowType) { this.targetRowType = targetRowType; } @Override public void open(InitializationContext context) { objectMapper = new ObjectMapper(); } @Override public Row deserialize(byte[] message) throws IOException { Map<String, Object> jsonMap = objectMapper.readValue(message, new TypeReference<>() {}); String[] fieldNames = targetRowType.getFieldNames(); DataType[] fieldTypes = targetRowType.getFieldTypes(); Object[] values = new Object[fieldNames.length]; for (int i = 0; i < fieldNames.length; i++) { String fieldName = fieldNames[i]; Object value = jsonMap.get(fieldName); // 类型转换适配 values[i] = convertValueToFlinkType(value, fieldTypes[i]); } return Row.of(values); } // 类型转换逻辑,需根据实际情况完善 private Object convertValueToFlinkType(Object value, DataType flinkType) { if (value == null) return null; if (flinkType instanceof IntType && value instanceof Number) { return ((Number) value).intValue(); } else if (flinkType instanceof BigIntType && value instanceof Number) { return ((Number) value).longValue(); } else if (flinkType instanceof DecimalType decimalType && value instanceof Number) { return BigDecimal.valueOf(((Number) value).doubleValue()).setScale(decimalType.getScale(), RoundingMode.HALF_UP); } else if (flinkType instanceof TimestampType && value instanceof String) { try { return Timestamp.valueOf((String) value); } catch (IllegalArgumentException e) { return null; } } return value; } @Override public boolean isEndOfStream(Row nextElement) { return false; } }
4. 动态写入Oracle
写入Oracle时,同样可以基于获取的表结构动态生成SQL语句,使用JdbcSink:
// 生成动态插入SQL String insertSql = generateInsertSql(targetRowType.getFieldNames(), params.get("oracle.table")); JdbcSink<RowData> jdbcSink = JdbcSink.sink( insertSql, (statement, rowData) -> { // 动态设置参数 for (int i = 0; i < targetRowType.getFieldCount(); i++) { Object value = rowData.getField(i); statement.setObject(i + 1, value); } }, JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchIntervalMs(5000) .build(), new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(params.get("oracle.jdbc.url")) .withUsername(params.get("oracle.username")) .withPassword(params.get("oracle.password")) .withDriverName("oracle.jdbc.driver.OracleDriver") .build() ); // 提交作业 env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "Kafka Source") .sinkTo(jdbcSink); env.execute("Kafka to Oracle Sync Job"); // 生成插入SQL的方法 private static String generateInsertSql(String[] fieldNames, String tableName) { String columns = String.join(", ", fieldNames); String placeholders = String.join(", ", Collections.nCopies(fieldNames.length, "?")); return "INSERT INTO " + tableName + " (" + columns + ") VALUES (" + placeholders + ")"; }
关键注意事项
- 类型映射兼容性:Oracle和Flink的类型映射需要覆盖业务场景,比如NUMBER类型要区分整数和小数,DATE/TIMESTAMP要统一格式。
- 容错机制:开启
ignoreParseErrors和failOnMissingField,避免单个错误数据导致作业失败,也可以结合SideOutput收集错误数据。 - 参数传递:所有配置(Kafka地址、Topic、Oracle连接信息、表名)通过
ParameterTool从命令行或配置文件传入,保证作业通用性。
内容的提问来源于stack exchange,提问作者Muhammed KESİCİ
相关产品推荐
相关产品推荐

