2929import java .sql .SQLException ;
3030import java .sql .Types ;
3131import java .util .ArrayList ;
32+ import java .util .HashMap ;
3233import java .util .List ;
34+ import java .util .Map ;
3335import java .util .Set ;
3436import 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