Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
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