Skip to content

Commit af67d36

Browse files
Improve STRUCT support and expand data type mappings
1 parent 5ab3cd8 commit af67d36

4 files changed

Lines changed: 92 additions & 23 deletions

File tree

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

Lines changed: 42 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -407,7 +407,21 @@ private void handleOracleSpecificType(ResultSet resultSet, StructuredRecord.Buil
407407

408408
private StructuredRecord convertStructToRecord(Struct struct, Schema schema, ResultSet resultSet)
409409
throws SQLException {
410-
Object[] attributes = struct.getAttributes();
410+
Object[] attributes;
411+
String sqlTypeName = "UNKNOWN";
412+
413+
try {
414+
sqlTypeName = struct.getSQLTypeName();
415+
if (sqlTypeName == null || sqlTypeName.trim().isEmpty()) {
416+
throw new InvalidStageException("Oracle Struct type name is missing or invalid.");
417+
}
418+
attributes = struct.getAttributes();
419+
} catch (SQLException e) {
420+
throw new InvalidStageException(
421+
String.format("Failed to retrieve attributes for Oracle Struct type '%s'. " +
422+
"Ensure the database connection is open and the type exists.", sqlTypeName), e);
423+
}
424+
411425
List<Schema.Field> fields = schema.getFields();
412426
StructuredRecord.Builder builder = StructuredRecord.builder(schema);
413427

@@ -421,16 +435,36 @@ private StructuredRecord convertStructToRecord(Struct struct, Schema schema, Res
421435
}
422436
// If it is an internal nested STRUCT, recurse down
423437
if (attrValue instanceof Struct) {
424-
Schema fieldSchema = field.getSchema().isNullable() ? field.getSchema().getNonNullable() : field.getSchema();
425-
builder.set(field.getName(), convertStructToRecord((Struct) attrValue, fieldSchema, resultSet));
438+
Struct nestedStruct = (Struct) attrValue;
439+
try {
440+
Schema fieldSchema = field.getSchema().isNullable() ? field.getSchema().getNonNullable() : field.getSchema();
441+
builder.set(field.getName(), convertStructToRecord(nestedStruct, fieldSchema, resultSet));
442+
} catch (Exception e) {
443+
String nestedTypeName = "UNKNOWN";
444+
try {
445+
nestedTypeName = nestedStruct.getSQLTypeName();
446+
} catch (SQLException ignored) {
447+
// Ignore if we can't fetch the type name during an ongoing failure
448+
}
449+
450+
throw new InvalidStageException(
451+
String.format("Failed to recursively process nested struct for field '%s' (Oracle Type: %s). " +
452+
"Check if the inner struct schema maps correctly.",
453+
field.getName(), nestedTypeName), e);
454+
}
426455
continue;
427456
}
428457

429-
String attrClassName = attrValue.getClass().getName();
430-
Schema fieldSchema = field.getSchema().isNullable() ? field.getSchema().getNonNullable() : field.getSchema();
431-
432-
OracleStructAttributeConverters.convertValue(builder, field, fieldSchema, attrValue, attrClassName,
433-
this::getBfileBytes);
458+
try {
459+
String attrClassName = attrValue.getClass().getName();
460+
Schema fieldSchema = field.getSchema().isNullable() ? field.getSchema().getNonNullable() : field.getSchema();
461+
OracleStructAttributeConverters.convertValue(builder, field, fieldSchema, attrValue, attrClassName,
462+
this::getBfileBytes);
463+
} catch (Exception e) {
464+
throw new InvalidStageException(
465+
String.format("Error converting flat attribute for field '%s' in struct '%s'.",
466+
field.getName(), sqlTypeName), e);
467+
}
434468
}
435469
return builder.build();
436470
}

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

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919
import com.google.common.collect.ImmutableSet;
2020
import io.cdap.cdap.api.data.schema.Schema;
2121
import io.cdap.plugin.db.CommonSchemaReader;
22-
import org.jetbrains.annotations.NotNull;
2322
import org.slf4j.Logger;
2423
import org.slf4j.LoggerFactory;
2524

@@ -215,11 +214,11 @@ private Schema getStructSchema(Connection connection, String typeName, String ow
215214
}
216215

217216
private Schema mapPrimitiveOracleType(String typeName, int precision, int scale, String columnName) {
218-
return OracleUserTypeSchemaMapping.mapPrimitiveOracleType(isTimestampOldBehavior, getTimestampLtzSchema(),
217+
return OracleStructTypeSchemaMapping.mapPrimitiveOracleType(isTimestampOldBehavior, getTimestampLtzSchema(),
219218
isPrecisionlessNumAsDecimal, typeName, precision, scale, columnName);
220219
}
221220

222-
private @NotNull Schema getTimestampLtzSchema() {
221+
private Schema getTimestampLtzSchema() {
223222
return isTimestampOldBehavior || isTimestampLtzFieldTimestamp
224223
? Schema.of(Schema.LogicalType.TIMESTAMP_MICROS)
225224
: Schema.of(Schema.LogicalType.DATETIME);

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

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import java.sql.Blob;
2424
import java.sql.Clob;
2525
import java.sql.SQLException;
26+
import java.sql.SQLXML;
2627
import java.sql.Timestamp;
2728
import java.time.OffsetDateTime;
2829
import java.time.ZoneId;
@@ -154,7 +155,13 @@ public void convert(StructuredRecord.Builder builder, Schema.Field field, Schema
154155
private static class OracleBfileConverter implements AttributeConverter {
155156
@Override
156157
public boolean canConvert(Object attrValue, String attrClassName) {
157-
return "oracle.jdbc.OracleBfile".equals(attrClassName);
158+
try {
159+
ClassLoader oracleLoader = attrValue.getClass().getClassLoader();
160+
Class<?> bfileInterface = oracleLoader.loadClass("oracle.jdbc.OracleBfile");
161+
return bfileInterface.isInstance(attrValue);
162+
} catch (Exception e) {
163+
return false;
164+
}
158165
}
159166

160167
@Override
@@ -190,6 +197,26 @@ public void convert(StructuredRecord.Builder builder, Schema.Field field, Schema
190197
}
191198
}
192199

200+
private static class SqlXmlConverter implements AttributeConverter {
201+
202+
@Override
203+
public boolean canConvert(Object attrValue, String attrClassName) {
204+
return attrValue instanceof SQLXML;
205+
}
206+
207+
@Override
208+
public void convert(
209+
StructuredRecord.Builder builder,
210+
Schema.Field field,
211+
Schema fieldSchema,
212+
Object attrValue,
213+
BfileBytesResolver resolver) throws SQLException {
214+
215+
SQLXML xml = (SQLXML) attrValue;
216+
builder.set(field.getName(), xml.getString());
217+
}
218+
}
219+
193220
private static class DefaultConverter implements AttributeConverter {
194221
@Override
195222
public boolean canConvert(Object attrValue, String attrClassName) {
@@ -212,6 +239,7 @@ public void convert(StructuredRecord.Builder builder, Schema.Field field, Schema
212239
new OracleBfileConverter(),
213240
new ByteArrayConverter(),
214241
new OracleIntervalConverter(),
242+
new SqlXmlConverter(),
215243
new DefaultConverter()
216244
);
217245

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

Lines changed: 19 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,8 @@
2929
/**
3030
* Registry containing schema type mappers for Oracle specific datatypes.
3131
*/
32-
public final class OracleUserTypeSchemaMapping {
33-
private static final Logger LOG = LoggerFactory.getLogger(OracleUserTypeSchemaMapping.class);
32+
public final class OracleStructTypeSchemaMapping {
33+
private static final Logger LOG = LoggerFactory.getLogger(OracleStructTypeSchemaMapping.class);
3434

3535
private interface TypeMapper {
3636
Schema map(boolean isTimestampOldBehavior, Schema timestampLtzSchema,
@@ -41,12 +41,12 @@ Schema map(boolean isTimestampOldBehavior, Schema timestampLtzSchema,
4141

4242
static {
4343
TypeMapper floatMapper = (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.Type.FLOAT);
44-
TYPE_MAPPERS.put("BINARY FLOAT", floatMapper);
44+
TYPE_MAPPERS.put("BINARY_FLOAT", floatMapper);
4545
TYPE_MAPPERS.put("REAL", floatMapper);
4646
TYPE_MAPPERS.put("FLOAT", floatMapper);
4747

4848
TypeMapper doubleMapper = (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.Type.DOUBLE);
49-
TYPE_MAPPERS.put("BINARY DOUBLE", doubleMapper);
49+
TYPE_MAPPERS.put("BINARY_DOUBLE", doubleMapper);
5050
TYPE_MAPPERS.put("DOUBLE", doubleMapper);
5151

5252
// Bytes types
@@ -64,44 +64,52 @@ Schema map(boolean isTimestampOldBehavior, Schema timestampLtzSchema,
6464
TYPE_MAPPERS.put("VARCHAR", stringMapper);
6565
TYPE_MAPPERS.put("CHAR", stringMapper);
6666
TYPE_MAPPERS.put("CHAR2", stringMapper);
67+
TYPE_MAPPERS.put("NCHAR", stringMapper);
68+
TYPE_MAPPERS.put("NVARCHAR2", stringMapper);
6769
TYPE_MAPPERS.put("CLOB", stringMapper);
6870
TYPE_MAPPERS.put("NCLOB", stringMapper);
6971
TYPE_MAPPERS.put("LONG", stringMapper);
72+
TYPE_MAPPERS.put("ROWID", stringMapper);
73+
TYPE_MAPPERS.put("UROWID", stringMapper);
7074

71-
// Specific types
75+
// Date and Time types
7276
TYPE_MAPPERS.put("TIMESTAMP WITH TZ", (isOld, ltzS, precD, typeName, p, s, col) ->
7377
isOld ? Schema.of(Schema.Type.STRING) : Schema.of(Schema.LogicalType.TIMESTAMP_MICROS)
7478
);
75-
TYPE_MAPPERS.put("TIMESTAMP WITH LTZ", (isOld, ltzS, precD, typeName, p, s, col) -> ltzS);
79+
TYPE_MAPPERS.put("TIMESTAMP WITH LOCAL TZ", (isOld, ltzS, precD, typeName, p, s, col) -> ltzS);
7680
TYPE_MAPPERS.put("TIMESTAMP", (isOld, ltzS, precD, typeName, p, s, col) ->
7781
isOld ? Schema.of(Schema.LogicalType.TIMESTAMP_MICROS) : Schema.of(Schema.LogicalType.DATETIME)
7882
);
7983
TYPE_MAPPERS.put("DATE", (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.LogicalType.DATE));
8084
TYPE_MAPPERS.put("TIME", (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.LogicalType.TIME_MICROS));
85+
86+
// Numeric types
8187
TYPE_MAPPERS.put("INTEGER", (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.Type.INT));
88+
TYPE_MAPPERS.put("NUMBER", OracleStructTypeSchemaMapping::mapNumberOrDecimal);
89+
TYPE_MAPPERS.put("DECIMAL", OracleStructTypeSchemaMapping::mapNumberOrDecimal);
8290

83-
TYPE_MAPPERS.put("NUMBER", OracleUserTypeSchemaMapping::mapNumberOrDecimal);
84-
TYPE_MAPPERS.put("DECIMAL", OracleUserTypeSchemaMapping::mapNumberOrDecimal);
91+
// XML type
92+
TYPE_MAPPERS.put("XMLTYPE", (isOld, ltzS, precD, typeName, p, s, col) -> Schema.of(Schema.Type.STRING));
8593

8694
// Unsupported types that throw error
8795
TYPE_MAPPERS.put("ARRAY", (isOld, ltzS, precD, typeName, p, s, col) -> {
8896
String errorMessage = String.format("Column %s has unsupported SQL type of %s.", col, typeName);
8997
throw ErrorUtils.getProgramFailureException(new ErrorCategory(ErrorCategory.ErrorCategoryEnum.PLUGIN),
9098
errorMessage, errorMessage, ErrorType.SYSTEM, true, null);
9199
});
92-
TYPE_MAPPERS.put("OTHER", (isOld, ltzS, precD, typeName, p, s, col) -> {
100+
TYPE_MAPPERS.put("ANYDATA", (isOld, ltzS, precD, typeName, p, s, col) -> {
93101
String errorMessage = String.format("Column %s has unsupported SQL type of %s.", col, typeName);
94102
throw ErrorUtils.getProgramFailureException(new ErrorCategory(ErrorCategory.ErrorCategoryEnum.PLUGIN),
95103
errorMessage, errorMessage, ErrorType.SYSTEM, true, null);
96104
});
97-
TYPE_MAPPERS.put("XML", (isOld, ltzS, precD, typeName, p, s, col) -> {
105+
TYPE_MAPPERS.put("OTHER", (isOld, ltzS, precD, typeName, p, s, col) -> {
98106
String errorMessage = String.format("Column %s has unsupported SQL type of %s.", col, typeName);
99107
throw ErrorUtils.getProgramFailureException(new ErrorCategory(ErrorCategory.ErrorCategoryEnum.PLUGIN),
100108
errorMessage, errorMessage, ErrorType.SYSTEM, true, null);
101109
});
102110
}
103111

104-
private OracleUserTypeSchemaMapping() {
112+
private OracleStructTypeSchemaMapping() {
105113
// Private constructor to prevent instantiation of utility class.
106114
}
107115

0 commit comments

Comments
 (0)