Skip to content

Commit 47834bb

Browse files
Updated schema mapping logic
1 parent 3bfd8ba commit 47834bb

7 files changed

Lines changed: 172 additions & 213 deletions

File tree

database-commons/src/main/java/io/cdap/plugin/db/CommonSchemaReader.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,11 @@ public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLExcepti
5757
metadata.isSigned(index), true);
5858
}
5959

60+
public Schema getSchema(String typeName, int sqlType, int precision, int scale, String columnName,
61+
boolean isSigned) throws SQLException {
62+
return DBUtils.getSchema(typeName, sqlType, precision, scale, columnName, isSigned, true);
63+
}
64+
6065
@Override
6166
public boolean shouldIgnoreColumn(ResultSetMetaData metadata, int index) throws SQLException {
6267
return false;

oracle-plugin/src/main/java/io/cdap/plugin/oracle/OracleSourceDBRecord.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -260,7 +260,7 @@ private byte[] getBfileBytes(ResultSet resultSet, String columnName) throws SQLE
260260
}
261261
}
262262

263-
private byte[] getBfileBytes(Object bfile) throws SQLException {
263+
static byte[] getBfileBytes(Object bfile) throws SQLException {
264264
if (bfile == null) {
265265
return null;
266266
}
@@ -429,8 +429,7 @@ private StructuredRecord convertStructToRecord(Struct struct, Schema schema, Res
429429
String attrClassName = attrValue.getClass().getName();
430430
Schema fieldSchema = field.getSchema().isNullable() ? field.getSchema().getNonNullable() : field.getSchema();
431431

432-
OracleStructAttributeConverters.convertValue(builder, field, fieldSchema, attrValue, attrClassName,
433-
this::getBfileBytes);
432+
OracleStructAttributeConverters.convertValue(builder, field, fieldSchema, attrValue, attrClassName);
434433
}
435434
return builder.build();
436435
}

oracle-plugin/src/main/java/io/cdap/plugin/oracle/OracleSourceSchemaReader.java

Lines changed: 84 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -29,7 +29,9 @@
2929
import java.sql.SQLException;
3030
import java.sql.Types;
3131
import java.util.ArrayList;
32+
import java.util.HashMap;
3233
import java.util.List;
34+
import java.util.Map;
3335
import java.util.Set;
3436
import javax.annotation.Nullable;
3537

@@ -50,6 +52,47 @@ public class OracleSourceSchemaReader extends CommonSchemaReader {
5052
public static final int LONG = -1;
5153
public static final int LONG_RAW = -4;
5254

55+
/**
56+
* Maps Oracle string data type inside UDT to their corresponding java.sql.Types integer constants
57+
*/
58+
private static final Map<String, Integer> DATA_TYPE_MAP = new HashMap<>();
59+
static {
60+
DATA_TYPE_MAP.put("TIMESTAMP WITH LOCAL TZ", TIMESTAMP_LTZ);
61+
DATA_TYPE_MAP.put("TIMESTAMP WITH TZ", TIMESTAMP_TZ);
62+
DATA_TYPE_MAP.put("TIMESTAMP", Types.TIMESTAMP);
63+
DATA_TYPE_MAP.put("DATE", Types.DATE);
64+
DATA_TYPE_MAP.put("TIME", Types.TIME);
65+
DATA_TYPE_MAP.put("FLOAT", Types.FLOAT);
66+
DATA_TYPE_MAP.put("BINARY_FLOAT", BINARY_FLOAT);
67+
DATA_TYPE_MAP.put("REAL", Types.REAL);
68+
DATA_TYPE_MAP.put("BINARY_DOUBLE", BINARY_DOUBLE);
69+
DATA_TYPE_MAP.put("DOUBLE", Types.DOUBLE);
70+
DATA_TYPE_MAP.put("BFILE", BFILE);
71+
DATA_TYPE_MAP.put("RAW", LONG_RAW);
72+
DATA_TYPE_MAP.put("LONG RAW", LONG_RAW);
73+
DATA_TYPE_MAP.put("LONG", LONG);
74+
DATA_TYPE_MAP.put("INTERVAL DAY TO SECOND", INTERVAL_DS);
75+
DATA_TYPE_MAP.put("INTERVAL YEAR TO MONTH", INTERVAL_YM);
76+
DATA_TYPE_MAP.put("XMLTYPE", Types.SQLXML);
77+
DATA_TYPE_MAP.put("ARRAY", Types.ARRAY);
78+
DATA_TYPE_MAP.put("ANYDATA", Types.JAVA_OBJECT);
79+
DATA_TYPE_MAP.put("OTHER", Types.OTHER);
80+
DATA_TYPE_MAP.put("NUMBER", Types.NUMERIC);
81+
DATA_TYPE_MAP.put("DECIMAL", Types.DECIMAL);
82+
DATA_TYPE_MAP.put("INTEGER", Types.INTEGER);
83+
DATA_TYPE_MAP.put("ROWID", Types.ROWID);
84+
DATA_TYPE_MAP.put("UROWID", Types.ROWID);
85+
DATA_TYPE_MAP.put("BLOB", Types.BLOB);
86+
DATA_TYPE_MAP.put("CLOB", Types.CLOB);
87+
DATA_TYPE_MAP.put("NCLOB", Types.NCLOB);
88+
DATA_TYPE_MAP.put("VARCHAR2", Types.VARCHAR);
89+
DATA_TYPE_MAP.put("VARCHAR", Types.VARCHAR);
90+
DATA_TYPE_MAP.put("CHAR", Types.CHAR);
91+
DATA_TYPE_MAP.put("CHAR2", Types.CHAR);
92+
DATA_TYPE_MAP.put("NCHAR", Types.NCHAR);
93+
DATA_TYPE_MAP.put("NVARCHAR2", Types.NVARCHAR);
94+
}
95+
5396
/**
5497
* Logger instance for Oracle Schema reader.
5598
*/
@@ -94,14 +137,27 @@ public OracleSourceSchemaReader(@Nullable String sessionID, boolean isTimestampO
94137
@Override
95138
public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLException {
96139
int sqlType = metadata.getColumnType(index);
140+
String owner = (metadata.getColumnTypeName(index) != null
141+
&& metadata.getColumnTypeName(index).contains(".")) ? metadata.getColumnTypeName(index)
142+
.substring(0, metadata.getColumnTypeName(index).lastIndexOf('.')) : null;
143+
144+
return getSchemaMapping(sqlType, metadata.getColumnClassName(index), metadata.getPrecision(index),
145+
metadata.getScale(index), metadata.getColumnName(index), metadata.getColumnTypeName(index),
146+
metadata.isSigned(index), owner, 0);
147+
}
148+
149+
public Schema getSchemaMapping(int sqlType, String columnClassName, int columnPrecision,
150+
int columnScale, String columnName, String columnTypeName,
151+
boolean isSigned, String owner, int nestingLevel) throws SQLException {
97152

98153
switch (sqlType) {
99154
case TIMESTAMP_TZ:
100155
return isTimestampOldBehavior ? Schema.of(Schema.Type.STRING) : Schema.of(Schema.LogicalType.TIMESTAMP_MICROS);
101156
case TIMESTAMP_LTZ:
102157
return getTimestampLtzSchema();
103158
case Types.TIMESTAMP:
104-
return isTimestampOldBehavior ? super.getSchema(metadata, index) : Schema.of(Schema.LogicalType.DATETIME);
159+
return isTimestampOldBehavior ? super.getSchema(columnTypeName, sqlType,
160+
columnPrecision, columnScale, columnName, isSigned) : Schema.of(Schema.LogicalType.DATETIME);
105161
case BINARY_FLOAT:
106162
return Schema.of(Schema.Type.FLOAT);
107163
case BINARY_DOUBLE:
@@ -115,15 +171,16 @@ public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLExcepti
115171
return Schema.of(Schema.Type.STRING);
116172
case Types.SQLXML:
117173
// Enabling XML type support for DTS connectors only as it is not in working state in CDAP plugin.
118-
return isXmlTypeEnabled ? Schema.of(Schema.Type.STRING) : super.getSchema(metadata, index);
174+
return isXmlTypeEnabled ? Schema.of(Schema.Type.STRING) : super.getSchema(columnTypeName,
175+
sqlType, columnPrecision, columnScale, columnName, isSigned);
119176
case Types.NUMERIC:
120177
case Types.DECIMAL:
121178
// FLOAT and REAL are returned as java.sql.Types.NUMERIC but with value that is a java.lang.Double
122-
if (Double.class.getTypeName().equals(metadata.getColumnClassName(index))) {
179+
if (Double.class.getTypeName().equals(columnClassName)) {
123180
return Schema.of(Schema.Type.DOUBLE);
124181
} else {
125-
int precision = metadata.getPrecision(index); // total number of digits
126-
int scale = metadata.getScale(index); // digits after the decimal point
182+
int precision = columnPrecision; // total number of digits
183+
int scale = columnScale; // digits after the decimal point
127184
// For a Number type without specified precision and scale, precision will be 0 and scale will be -127
128185
if (precision == 0) {
129186
// reference : https://docs.oracle.com/cd/B28359_01/server.111/b28318/datatype.htm#CNCPT1832
@@ -134,15 +191,12 @@ public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLExcepti
134191
+ "there may be a precision loss while running the pipeline. "
135192
+ "Please define an output precision and scale for field '%s' to avoid "
136193
+ "precision loss.",
137-
metadata.getColumnTypeName(index),
138-
metadata.getColumnName(index)));
194+
columnTypeName, columnName));
139195
return Schema.decimalOf(precision, scale);
140196
} else {
141197
LOG.warn(String.format("Field '%s' is a %s type without precision and scale, "
142198
+ "converting into STRING type to avoid any precision loss.",
143-
metadata.getColumnName(index),
144-
metadata.getColumnTypeName(index),
145-
metadata.getColumnName(index)));
199+
columnName, columnTypeName, columnName));
146200
return Schema.of(Schema.Type.STRING);
147201
}
148202
}
@@ -153,11 +207,13 @@ public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLExcepti
153207
throw new SQLException("Cannot resolve STRUCT schema without a database connection. "
154208
+ "Use getSchemaFields(ResultSet) to enable STRUCT type resolution.");
155209
}
156-
String typeName = metadata.getColumnTypeName(index);
157-
String owner = typeName.substring(0, typeName.lastIndexOf('.'));
158-
return getStructSchema(connection, typeName, owner);
210+
if (nestingLevel >= 4) {
211+
throw new IllegalArgumentException(String.format("Cannot resolve STRUCT schema for attribute %s with " +
212+
"nested structure depth more than 4.", columnName));
213+
}
214+
return getStructSchema(connection, columnTypeName, owner, nestingLevel);
159215
default:
160-
return super.getSchema(metadata, index);
216+
return super.getSchema(columnTypeName, sqlType, columnPrecision, columnScale, columnName, isSigned);
161217
}
162218
}
163219

@@ -174,10 +230,12 @@ public List<Schema.Field> getSchemaFields(ResultSet resultSet) throws SQLExcepti
174230
*
175231
* @param connection the database connection
176232
* @param typeName the Oracle type name (e.g., "ADDRESS_TYPE")
233+
* @param owner the Owner of the user-defined data type
234+
* @param level the level of nesting of the user-defined data type
177235
* @return a CDAP RECORD schema with fields corresponding to the STRUCT's
178236
* attributes
179237
*/
180-
private Schema getStructSchema(Connection connection, String typeName, String owner) throws SQLException {
238+
private Schema getStructSchema(Connection connection, String typeName, String owner, int level) throws SQLException {
181239
List<Schema.Field> fields = new ArrayList<>();
182240
String sql = "SELECT * FROM ALL_TYPE_ATTRS WHERE TYPE_NAME = ? AND OWNER = ? ORDER BY ATTR_NO";
183241

@@ -191,22 +249,25 @@ private Schema getStructSchema(Connection connection, String typeName, String ow
191249
String attrTypeName = attrRs.getString("ATTR_TYPE_NAME");
192250
int attrSize = attrRs.getInt("PRECISION");
193251
int attrScale = attrRs.getInt("SCALE");
252+
Integer sqlType = DATA_TYPE_MAP.getOrDefault(attrTypeName, null);
194253

195-
Schema attrSchema = mapPrimitiveOracleType(attrTypeName, attrSize, attrScale, attrName);
196-
if (attrSchema != null) {
197-
fields.add(Schema.Field.of(attrName, attrSchema));
198-
} else {
199-
String nestedStructOwner = attrRs.getString("ATTR_TYPE_OWNER");
200-
if (nestedStructOwner == null || nestedStructOwner.isEmpty()) {
254+
int nextLevel = level;
255+
if (sqlType == null) {
256+
owner = attrRs.getString("ATTR_TYPE_OWNER");
257+
if (owner == null || owner.isEmpty()) {
201258
throw new SQLException(String.format("Attribute '%s' is not a primitive type, but it lacks a type " +
202259
"owner. Therefore, it cannot be resolved as a STRUCT type. ", attrName));
203260
}
204-
Schema nestedSchema = getStructSchema(connection, attrTypeName, nestedStructOwner);
205-
fields.add(Schema.Field.of(attrName, nestedSchema));
261+
sqlType = Types.STRUCT;
262+
nextLevel = level + 1;
206263
}
264+
Schema attrSchema = getSchemaMapping(sqlType, null, attrSize,
265+
attrScale, attrName, attrTypeName, true, owner, nextLevel);
266+
fields.add(Schema.Field.of(attrName, attrSchema));
207267
}
208268
}
209269
}
270+
210271
if (fields.isEmpty()) {
211272
throw new SQLException(String.format(
212273
"No attributes found for Oracle STRUCT type '%s'. "
@@ -217,11 +278,6 @@ private Schema getStructSchema(Connection connection, String typeName, String ow
217278
return Schema.recordOf(typeName, fields);
218279
}
219280

220-
private Schema mapPrimitiveOracleType(String typeName, int precision, int scale, String columnName) {
221-
return OracleStructTypeSchemaMapping.mapPrimitiveOracleType(isTimestampOldBehavior, getTimestampLtzSchema(),
222-
isPrecisionlessNumAsDecimal, typeName, precision, scale, columnName);
223-
}
224-
225281
private Schema getTimestampLtzSchema() {
226282
return isTimestampOldBehavior || isTimestampLtzFieldTimestamp
227283
? Schema.of(Schema.LogicalType.TIMESTAMP_MICROS)

0 commit comments

Comments
 (0)