Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,11 @@ public Schema getSchema(ResultSetMetaData metadata, int index) throws SQLExcepti
metadata.isSigned(index), true);
}

public Schema getSchema(String typeName, int sqlType, int precision, int scale, String columnName,
boolean isSigned) throws SQLException {
return DBUtils.getSchema(typeName, sqlType, precision, scale, columnName, isSigned, true);
}

@Override
public boolean shouldIgnoreColumn(ResultSetMetaData metadata, int index) throws SQLException {
return false;
Expand Down
18 changes: 17 additions & 1 deletion database-commons/src/main/java/io/cdap/plugin/db/DBRecord.java
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import java.math.BigDecimal;
import java.math.BigInteger;
import java.nio.ByteBuffer;
import java.sql.Connection;
import java.sql.Date;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
Expand Down Expand Up @@ -187,7 +188,22 @@ protected void handleField(ResultSet resultSet, StructuredRecord.Builder recordB

protected void setField(ResultSet resultSet, StructuredRecord.Builder recordBuilder, Schema.Field field,
int columnIndex, int sqlType, int sqlPrecision, int sqlScale) throws SQLException {
Object o = DBUtils.transformValue(sqlType, sqlPrecision, sqlScale, resultSet, columnIndex);
Object fieldValue = DBUtils.transformValue(sqlType, sqlPrecision, sqlScale, resultSet, columnIndex);
populateRecordField(resultSet.getStatement().getConnection(), recordBuilder, field, fieldValue);
}

/**
* Populates the value of a field in the {@link StructuredRecord.Builder}.
*
* @param connection the SQL connection, provided for subclass overrides that require database connection
* @param recordBuilder the builder for constructing the {@link StructuredRecord}
* @param field the field to set in the record
* @param o the object value read from the database
* @throws SQLException if an error occurs while setting the field value
*/
public void populateRecordField(Connection connection, StructuredRecord.Builder recordBuilder,
Schema.Field field, Object o)
throws SQLException {
if (o instanceof Date) {
recordBuilder.setDate(field.getName(), ((Date) o).toLocalDate());
} else if (o instanceof Time) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
import java.sql.ResultSet;
import java.sql.ResultSetMetaData;
import java.sql.SQLException;
import java.sql.Struct;
import java.sql.Timestamp;
import java.sql.Types;
import java.time.LocalDateTime;
Expand All @@ -43,6 +44,8 @@
import java.time.ZoneOffset;
import java.time.ZonedDateTime;
import java.util.List;
import java.util.Map;
import java.util.TreeMap;

/**
* Oracle Source implementation {@link org.apache.hadoop.mapreduce.lib.db.DBWritable} and
Expand Down Expand Up @@ -106,13 +109,19 @@ record = recordBuilder.build();
@Override
protected void handleField(ResultSet resultSet, StructuredRecord.Builder recordBuilder, Schema.Field field,
int columnIndex, int sqlType, int sqlPrecision, int sqlScale) throws SQLException {
if (OracleSourceSchemaReader.ORACLE_TYPES.contains(sqlType) || sqlType == Types.NCLOB) {
if (isOracleSpecificType(sqlType)) {
handleOracleSpecificType(resultSet, recordBuilder, field, columnIndex, sqlType, sqlPrecision, sqlScale);
} else {
setField(resultSet, recordBuilder, field, columnIndex, sqlType, sqlPrecision, sqlScale);
}
}

protected boolean isOracleSpecificType(int sqlType) {
return OracleSourceSchemaReader.ORACLE_TYPES.contains(sqlType)
|| sqlType == Types.NCLOB
|| sqlType == Types.STRUCT;
}

@Override
protected void writeNonNullToDB(PreparedStatement stmt, Schema fieldSchema,
String fieldName, int fieldIndex) throws SQLException {
Expand Down Expand Up @@ -232,11 +241,15 @@ private Object createOracleTimestamp(Connection connection, String timestampStri
*/
private byte[] getBfileBytes(ResultSet resultSet, String columnName) throws SQLException {
Object bfile = resultSet.getObject(columnName);
return getBfileBytes(bfile, columnName);
}

public byte[] getBfileBytes(Object bfile, String columnName) {
if (bfile == null) {
return null;
}
try {
ClassLoader classLoader = resultSet.getClass().getClassLoader();
ClassLoader classLoader = bfile.getClass().getClassLoader();
Class<?> oracleBfileClass = classLoader.loadClass("oracle.jdbc.OracleBfile");
boolean isFileExist = (boolean) oracleBfileClass.getMethod("fileExists").invoke(bfile);
if (!isFileExist) {
Expand Down Expand Up @@ -341,6 +354,15 @@ private void handleOracleSpecificType(ResultSet resultSet, StructuredRecord.Buil
case OracleSourceSchemaReader.LONG_RAW:
recordBuilder.set(field.getName(), resultSet.getBytes(columnIndex));
break;
case Types.STRUCT:
Struct structValue = (Struct) resultSet.getObject(columnIndex);
if (structValue != null) {
recordBuilder.set(field.getName(), convertStructToRecord(structValue, nonNullSchema,
resultSet.getStatement().getConnection()));
} else {
recordBuilder.set(field.getName(), null);
}
break;
case Types.DECIMAL:
case Types.NUMERIC:
// This is the only way to differentiate FLOAT/REAL columns from other numeric columns, that based on NUMBER.
Expand Down Expand Up @@ -371,6 +393,57 @@ private void handleOracleSpecificType(ResultSet resultSet, StructuredRecord.Buil
}
}

/**
* Converts a JDBC {@link Struct} into a {@link StructuredRecord} based on the provided schema.
*
* @param struct the SQL structured type containing the source data attributes
* @param schema the target record schema defining the fields to map
* @param connection the database connection
* @return a populated {@code StructuredRecord} instance
* @throws SQLException if an error occurs reading the struct attributes or metadata
*/
protected StructuredRecord convertStructToRecord(Struct struct, Schema schema, Connection connection)
throws SQLException {
Map<String, Object> attributeMap = getAttributeMap(struct, schema, connection);
StructuredRecord.Builder builder = StructuredRecord.builder(schema);

for (Schema.Field field : schema.getFields()) {
Object attrValue = attributeMap.get(field.getName());
OracleStructUtil.populateRecordField(this, connection, builder, field, attrValue);
}
return builder.build();
}

/**
* Extracts attributes from a {@link Struct} into a case-insensitive map indexed by column name.
* Uses reflection to extract underlying metadata (e.g., from Oracle StructDescriptor).
*
* @param struct the source SQL structured type
* @param schema the target schema used for context in error messages
* @param connection the database connection
* @return a case-insensitive {@code Map} linking column names to their attribute values
* @throws SQLException if metadata extraction fails or driver-specific methods are inaccessible
*/
protected Map<String, Object> getAttributeMap(Struct struct, Schema schema, Connection connection)
throws SQLException {
Map<String, Object> attributeMap = new TreeMap<>(String.CASE_INSENSITIVE_ORDER);
Object[] attributes = struct.getAttributes();
if (attributes != null) {
try {
Object descriptor = struct.getClass().getMethod("getDescriptor").invoke(struct);
Comment thread
vanshikaagupta22 marked this conversation as resolved.
ResultSetMetaData metaData =
(ResultSetMetaData) descriptor.getClass().getMethod("getMetaData").invoke(descriptor);
for (int i = 1; i <= metaData.getColumnCount() && (i - 1) < attributes.length; i++) {
attributeMap.put(metaData.getColumnName(i), attributes[i - 1]);
}
} catch (Exception e) {
throw new SQLException(String.format("Failed to retrieve attribute metadata for Oracle STRUCT schema '%s': %s",
schema.getRecordName(), e.getMessage()));
}
}
return attributeMap;
}

/**
* Get the scale set in Non-nullable schema associated with the schema
* */
Expand Down
Loading
Loading