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

如何实现通用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İ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 21:55:58