From e7bbd100848f7b30624ae043feae12ceeb882f58 Mon Sep 17 00:00:00 2001 From: xxsc0529 Date: Tue, 7 Jul 2026 16:06:12 +0800 Subject: [PATCH] feat: oceanbase oracle mode support string type split --- .../oceanbasev10reader/ext/ReaderJob.java | 3 +- .../util/ObReaderSplitUtil.java | 92 +++++++ .../util/ObSingleTableSplitUtil.java | 249 ++++++++++++++++++ 3 files changed, 343 insertions(+), 1 deletion(-) create mode 100644 oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObReaderSplitUtil.java create mode 100644 oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObSingleTableSplitUtil.java diff --git a/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/ext/ReaderJob.java b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/ext/ReaderJob.java index 2d60d0c6..02070933 100644 --- a/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/ext/ReaderJob.java +++ b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/ext/ReaderJob.java @@ -8,6 +8,7 @@ import com.alibaba.datax.plugin.rdbms.reader.CommonRdbmsReader; import com.alibaba.datax.plugin.rdbms.reader.Key; import com.alibaba.datax.plugin.rdbms.reader.Constant; import com.alibaba.datax.plugin.reader.oceanbasev10reader.OceanBaseReader; +import com.alibaba.datax.plugin.reader.oceanbasev10reader.util.ObReaderSplitUtil; import com.alibaba.datax.plugin.reader.oceanbasev10reader.util.ObReaderUtils; import com.alibaba.datax.plugin.reader.oceanbasev10reader.util.PartitionSplitUtil; import com.alibaba.fastjson2.JSONObject; @@ -57,7 +58,7 @@ public class ReaderJob extends CommonRdbmsReader.Job { list = PartitionSplitUtil.splitByPartition(originalConfig); } else { LOG.info("try to split reader job by splitPk."); - list = super.split(originalConfig, adviceNumber); + list = ObReaderSplitUtil.doSplit(originalConfig, adviceNumber); } for (Configuration config : list) { diff --git a/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObReaderSplitUtil.java b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObReaderSplitUtil.java new file mode 100644 index 00000000..8721fa0d --- /dev/null +++ b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObReaderSplitUtil.java @@ -0,0 +1,92 @@ +package com.alibaba.datax.plugin.reader.oceanbasev10reader.util; + +import com.alibaba.datax.common.constant.CommonConstant; +import com.alibaba.datax.common.util.Configuration; +import com.alibaba.datax.plugin.rdbms.reader.Constant; +import com.alibaba.datax.plugin.rdbms.reader.Key; +import com.alibaba.datax.plugin.rdbms.reader.util.HintUtil; +import com.alibaba.datax.plugin.rdbms.reader.util.SingleTableSplitUtil; +import com.alibaba.datax.plugin.rdbms.util.DataBaseType; +import org.apache.commons.lang3.StringUtils; +import org.apache.commons.lang3.Validate; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.ArrayList; +import java.util.List; + +/** + * OceanBase Reader 专用的 Job 切分逻辑,在 splitPk 场景下使用 {@link ObSingleTableSplitUtil}。 + */ +public final class ObReaderSplitUtil { + private static final Logger LOG = LoggerFactory.getLogger(ObReaderSplitUtil.class); + + private ObReaderSplitUtil() { + } + + public static List doSplit(Configuration originalSliceConfig, int adviceNumber) { + boolean isTableMode = originalSliceConfig.getBool(Constant.IS_TABLE_MODE).booleanValue(); + int eachTableShouldSplittedNumber = -1; + if (isTableMode) { + eachTableShouldSplittedNumber = calculateEachTableShouldSplittedNumber( + adviceNumber, originalSliceConfig.getInt(Constant.TABLE_NUMBER_MARK)); + } + + String column = originalSliceConfig.getString(Key.COLUMN); + String where = originalSliceConfig.getString(Key.WHERE, null); + List conns = originalSliceConfig.getList(Constant.CONN_MARK, Object.class); + List splittedConfigs = new ArrayList(); + + for (int i = 0, len = conns.size(); i < len; i++) { + Configuration sliceConfig = originalSliceConfig.clone(); + Configuration connConf = Configuration.from(conns.get(i).toString()); + String jdbcUrl = connConf.getString(Key.JDBC_URL); + sliceConfig.set(Key.JDBC_URL, jdbcUrl); + sliceConfig.set(CommonConstant.LOAD_BALANCE_RESOURCE_MARK, DataBaseType.parseIpFromJdbcUrl(jdbcUrl)); + sliceConfig.remove(Constant.CONN_MARK); + + if (isTableMode) { + List tables = connConf.getList(Key.TABLE, String.class); + Validate.isTrue(null != tables && !tables.isEmpty(), "您读取数据库表配置错误."); + + String splitPk = originalSliceConfig.getString(Key.SPLIT_PK, null); + boolean needSplitTable = eachTableShouldSplittedNumber > 1 + && StringUtils.isNotBlank(splitPk); + if (needSplitTable) { + if (tables.size() == 1) { + Integer splitFactor = originalSliceConfig.getInt(Key.SPLIT_FACTOR, Constant.SPLIT_FACTOR); + eachTableShouldSplittedNumber = eachTableShouldSplittedNumber * splitFactor; + } + for (String table : tables) { + Configuration tempSlice = sliceConfig.clone(); + tempSlice.set(Key.TABLE, table); + splittedConfigs.addAll( + ObSingleTableSplitUtil.splitSingleTable(tempSlice, eachTableShouldSplittedNumber)); + } + } else { + for (String table : tables) { + Configuration tempSlice = sliceConfig.clone(); + tempSlice.set(Key.TABLE, table); + String queryColumn = HintUtil.buildQueryColumn(jdbcUrl, table, column); + tempSlice.set(Key.QUERY_SQL, + SingleTableSplitUtil.buildQuerySql(queryColumn, table, where)); + splittedConfigs.add(tempSlice); + } + } + } else { + List sqls = connConf.getList(Key.QUERY_SQL, String.class); + for (String querySql : sqls) { + Configuration tempSlice = sliceConfig.clone(); + tempSlice.set(Key.QUERY_SQL, querySql); + splittedConfigs.add(tempSlice); + } + } + } + + return splittedConfigs; + } + + private static int calculateEachTableShouldSplittedNumber(int adviceNumber, int tableNumber) { + return (int) Math.ceil(1.0 * adviceNumber / tableNumber); + } +} diff --git a/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObSingleTableSplitUtil.java b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObSingleTableSplitUtil.java new file mode 100644 index 00000000..e72812ab --- /dev/null +++ b/oceanbasev10reader/src/main/java/com/alibaba/datax/plugin/reader/oceanbasev10reader/util/ObSingleTableSplitUtil.java @@ -0,0 +1,249 @@ +package com.alibaba.datax.plugin.reader.oceanbasev10reader.util; + +import com.alibaba.datax.common.exception.DataXException; +import com.alibaba.datax.common.util.Configuration; +import com.alibaba.datax.plugin.rdbms.reader.Constant; +import com.alibaba.datax.plugin.rdbms.reader.Key; +import com.alibaba.datax.plugin.rdbms.reader.util.SingleTableSplitUtil; +import com.alibaba.datax.plugin.rdbms.util.DBUtil; +import com.alibaba.datax.plugin.rdbms.util.DBUtilErrorCode; +import com.alibaba.datax.plugin.rdbms.util.DataBaseType; +import com.alibaba.datax.plugin.rdbms.util.RdbmsException; +import com.alibaba.datax.plugin.rdbms.util.RdbmsRangeSplitWrap; +import com.alibaba.datax.plugin.reader.oceanbasev10reader.ext.ObReaderKey; +import org.apache.commons.lang3.StringUtils; +import org.apache.commons.lang3.tuple.ImmutablePair; +import org.apache.commons.lang3.tuple.Pair; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.math.BigInteger; +import java.sql.Connection; +import java.sql.ResultSet; +import java.sql.ResultSetMetaData; +import java.sql.Statement; +import java.sql.Types; +import java.util.ArrayList; +import java.util.List; + +/** + * OceanBase Reader 专用的单表切分逻辑,包含字符串 splitPk(Oracle 模式)及 OB V10 min/max 查询优化。 + */ +public final class ObSingleTableSplitUtil { + private static final Logger LOG = LoggerFactory.getLogger(ObSingleTableSplitUtil.class); + + private static final DataBaseType DATABASE_TYPE = ObReaderUtils.databaseType; + + private ObSingleTableSplitUtil() { + } + + public static List splitSingleTable(Configuration configuration, int adviceNum) { + List pluginParams = new ArrayList(); + List rangeList; + String splitPkName = configuration.getString(Key.SPLIT_PK); + String column = configuration.getString(Key.COLUMN); + String table = configuration.getString(Key.TABLE); + String where = configuration.getString(Key.WHERE, null); + boolean hasWhere = StringUtils.isNotBlank(where); + + Pair minMaxPK = getPkRange(configuration); + if (null == minMaxPK) { + throw DataXException.asDataXException(DBUtilErrorCode.ILLEGAL_SPLIT_PK, + "根据切分主键切分表失败. DataX 仅支持切分主键为一个,并且类型为整数或者字符串类型. 请尝试使用其他的切分主键或者联系 DBA 进行处理."); + } + + configuration.set(Key.QUERY_SQL, SingleTableSplitUtil.buildQuerySql(column, table, where)); + if (null == minMaxPK.getLeft() || null == minMaxPK.getRight()) { + pluginParams.add(configuration); + return pluginParams; + } + + boolean isStringType = Constant.PK_TYPE_STRING.equals(configuration.getString(Constant.PK_TYPE)); + boolean isLongType = Constant.PK_TYPE_LONG.equals(configuration.getString(Constant.PK_TYPE)); + + if (isStringType) { + DataBaseType stringSplitDbType = resolveStringSplitDataBaseType(configuration); + if (stringSplitDbType == null) { + pluginParams.add(configuration); + LOG.warn("切分主键 {} 为字符串类型,当前 OceanBase MySQL 模式不支持按字符串切分,将使用单通道读取。", + splitPkName); + return pluginParams; + } + rangeList = RdbmsRangeSplitWrap.splitAndWrap( + String.valueOf(minMaxPK.getLeft()), + String.valueOf(minMaxPK.getRight()), adviceNum, + splitPkName, "'", stringSplitDbType); + } else if (isLongType) { + rangeList = RdbmsRangeSplitWrap.splitAndWrap( + new BigInteger(minMaxPK.getLeft().toString()), + new BigInteger(minMaxPK.getRight().toString()), + adviceNum, splitPkName); + } else { + throw DataXException.asDataXException(DBUtilErrorCode.ILLEGAL_SPLIT_PK, + "您配置的切分主键(splitPk) 类型 DataX 不支持. DataX 仅支持切分主键为一个,并且类型为整数或者字符串类型. 请尝试使用其他的切分主键或者联系 DBA 进行处理."); + } + + String tempQuerySql; + List allQuerySql = new ArrayList(); + + if (null != rangeList && !rangeList.isEmpty()) { + for (String range : rangeList) { + Configuration tempConfig = configuration.clone(); + tempQuerySql = SingleTableSplitUtil.buildQuerySql(column, table, where) + + (hasWhere ? " and " : " where ") + range; + allQuerySql.add(tempQuerySql); + tempConfig.set(Key.QUERY_SQL, tempQuerySql); + tempConfig.set(Key.WHERE, (hasWhere ? ("(" + where + ") and") : "") + range); + pluginParams.add(tempConfig); + } + } else { + Configuration tempConfig = configuration.clone(); + tempQuerySql = SingleTableSplitUtil.buildQuerySql(column, table, where) + + (hasWhere ? " and " : " where ") + + String.format(" %s IS NOT NULL", splitPkName); + allQuerySql.add(tempQuerySql); + tempConfig.set(Key.QUERY_SQL, tempQuerySql); + tempConfig.set(Key.WHERE, (hasWhere ? "(" + where + ") and" : "") + + String.format(" %s IS NOT NULL", splitPkName)); + pluginParams.add(tempConfig); + } + + Configuration tempConfig = configuration.clone(); + tempQuerySql = SingleTableSplitUtil.buildQuerySql(column, table, where) + + (hasWhere ? " and " : " where ") + + String.format(" %s IS NULL", splitPkName); + allQuerySql.add(tempQuerySql); + LOG.info("After split(), allQuerySql=[\n{}\n].", StringUtils.join(allQuerySql, "\n")); + tempConfig.set(Key.QUERY_SQL, tempQuerySql); + tempConfig.set(Key.WHERE, (hasWhere ? "(" + where + ") and" : "") + + String.format(" %s IS NULL", splitPkName)); + pluginParams.add(tempConfig); + + return pluginParams; + } + + private static Pair getPkRange(Configuration configuration) { + int fetchSize = configuration.getInt(Constant.FETCH_SIZE); + String jdbcURL = configuration.getString(Key.JDBC_URL); + String username = configuration.getString(Key.USERNAME); + String password = configuration.getString(Key.PASSWORD); + String table = configuration.getString(Key.TABLE); + + Connection conn = DBUtil.getConnection(DATABASE_TYPE, jdbcURL, username, password); + Pair minMaxPK = checkSplitPk(conn, fetchSize, table, username, configuration); + DBUtil.closeDBResources(null, null, conn); + return minMaxPK; + } + + private static Pair checkSplitPk(Connection conn, int fetchSize, String table, + String username, Configuration configuration) { + LOG.info("Get min/max of split key for OBV10"); + int queryTimeoutSeconds = 60 * 60 * 48; + String setQueryTimeout = "set ob_query_timeout=" + (queryTimeoutSeconds * 1000 * 1000L); + String setTrxTimeout = "set ob_trx_timeout=" + ((queryTimeoutSeconds + 5) * 1000 * 1000L); + Statement stmt = null; + try { + stmt = conn.createStatement(); + stmt.execute(setQueryTimeout); + stmt.execute(setTrxTimeout); + } catch (Exception e) { + LOG.warn("set ob_query_timeout and set ob_trx_timeout failed. reason: {}", e.getMessage(), e); + } finally { + DBUtil.closeDBResources(stmt, null); + } + return new ImmutablePair( + getValueForObV10(conn, "min", configuration), + getValueForObV10(conn, "max", configuration)); + } + + private static Object getValueForObV10(Connection conn, String function, Configuration configuration) { + ResultSet rs = null; + Object value = null; + final int fetchSize = 1; + String username = configuration.getString(Key.USERNAME); + String table = configuration.getString(Key.TABLE); + String sql = genSqlForObV10(configuration, function); + + LOG.info("Running Query [{}]", sql); + + try { + rs = DBUtil.query(conn, sql, fetchSize); + ResultSetMetaData rsMetaData = rs.getMetaData(); + if (isPKTypeValid(rsMetaData)) { + if (isStringType(rsMetaData.getColumnType(1))) { + configuration.set(Constant.PK_TYPE, Constant.PK_TYPE_STRING); + } else if (isLongType(rsMetaData.getColumnType(1))) { + configuration.set(Constant.PK_TYPE, Constant.PK_TYPE_LONG); + } else { + throw DataXException.asDataXException(DBUtilErrorCode.ILLEGAL_SPLIT_PK, + "您配置的DataX切分主键(splitPk)有误. 因为您配置的切分主键(splitPk) 类型 DataX 不支持. DataX 仅支持切分主键为一个,并且类型为整数或者字符串类型. 请尝试使用其他的切分主键或者联系 DBA 进行处理."); + } + while (DBUtil.asyncResultSetNext(rs)) { + value = rs.getString(1); + } + } else { + throw DataXException.asDataXException(DBUtilErrorCode.ILLEGAL_SPLIT_PK, + "您配置的DataX切分主键(splitPk)有误. 因为您配置的切分主键(splitPk) 类型 DataX 不支持. DataX 仅支持切分主键为一个,并且类型为整数或者字符串类型. 请尝试使用其他的切分主键或者联系 DBA 进行处理."); + } + } catch (DataXException e) { + throw e; + } catch (Exception e) { + throw RdbmsException.asQueryException(DATABASE_TYPE, e, sql, table, username); + } finally { + DBUtil.closeDBResources(rs, null, null); + } + return value; + } + + private static String genSqlForObV10(Configuration configuration, String function) { + String primaryKey = configuration.getString(Key.SPLIT_PK).trim(); + String table = configuration.getString(Key.TABLE).trim(); + String where = configuration.getString(Key.WHERE, null); + String pkRangeSQL = String.format("SELECT %s(%s) FROM %s", function, primaryKey, table); + if (StringUtils.isNotBlank(where)) { + pkRangeSQL = String.format("%s WHERE (%s AND %s IS NOT NULL)", pkRangeSQL, where, primaryKey); + } + return pkRangeSQL; + } + + private static DataBaseType resolveStringSplitDataBaseType(Configuration configuration) { + if (isObOracleMode(configuration)) { + return DataBaseType.Oracle; + } + return null; + } + + private static boolean isObOracleMode(Configuration configuration) { + return configuration != null + && configuration.getString(ObReaderKey.OB_COMPATIBILITY_MODE, "") + .equalsIgnoreCase(DataBaseType.Oracle.getTypeName()); + } + + private static boolean isPKTypeValid(ResultSetMetaData rsMetaData) { + try { + int pkType = rsMetaData.getColumnType(1); + int maxType = Types.NULL; + if (rsMetaData.getColumnCount() == 2) { + maxType = rsMetaData.getColumnType(2); + } + boolean isNumberType = isLongType(pkType); + boolean isStringType = isStringType(pkType); + return (maxType == Types.NULL || pkType == maxType) && (isNumberType || isStringType); + } catch (Exception e) { + throw DataXException.asDataXException(DBUtilErrorCode.ILLEGAL_SPLIT_PK, + "DataX获取切分主键(splitPk)字段类型失败. 该错误通常是系统底层异常导致. 请联DBA进行处理."); + } + } + + private static boolean isLongType(int type) { + boolean isValidLongType = type == Types.BIGINT || type == Types.INTEGER + || type == Types.SMALLINT || type == Types.TINYINT || type == Types.NUMERIC; + return isValidLongType; + } + + private static boolean isStringType(int type) { + return type == Types.CHAR || type == Types.NCHAR + || type == Types.VARCHAR || type == Types.LONGVARCHAR + || type == Types.NVARCHAR; + } +}