使用parquet-avro 1.12.3时AvroParquetWriter的addLogicalTypeConversion失效引发ClassCastException
使用AvroParquetWriter将包含java.sql.Timestamp类型列的ResultSet写入Parquet文件时,抛出如下异常:
java.sql.Timestamp cannot be cast to java.lang.Number
已尝试添加addLogicalTypeConversion但未解决问题,相关代码如下,需要给出让LogicalTypeConversion生效的解决方案。
核心代码片段
写入逻辑
GenericData timeSupport = GenericData.get(); timeSupport.addLogicalTypeConversion(new TimeConversions.DateConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.LocalTimestampMillisConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.LocalTimestampMicrosConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.TimeMicrosConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.TimeMillisConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.TimestampMicrosConversion()); timeSupport.addLogicalTypeConversion(new TimeConversions.TimestampMillisConversion()); OutputFile opHadoop = HadoopOutputFile.fromPath(outputPath, new Configuration()); SchemaResults schemaResults = new ResultSetSchemaGenerator().generateSchema(resultSet, schemaName, namespace); ParquetWriter<GenericRecord> parquetWriter = AvroParquetWriter.<GenericRecord> builder(opHadoop) .withSchema(schemaResults.getParsedSchema()) .withDataModel(timeSupport) .withCompressionCodec(CompressionCodecName.SNAPPY) .build(); List<GenericRecord> records = new ArrayList<>(); while (resultSet.next()) { GenericRecordBuilder builder = new GenericRecordBuilder(schemaResults.getParsedSchema()) for (SchemaSqlMapping mapping : schemaResults.getMappings()) { builder.set( schemaResults.getParsedSchema().getField(mapping.getSchemaName()), extractResult(mapping, resultSet)); } GenericRecord record = builder.build(); records.add(record); } for (GenericRecord record : records) { parquetWriter.write(record); } parquetWriter.close();
结果集提取方法
/** * Extracts the appropriate value from the ResultSet using the given * SchemaSqlMapping. * * @param mapping * @param resultSet * @return * @throws SQLException */ public Object extractResult(SchemaSqlMapping mapping, ResultSet resultSet) throws SQLException { switch (mapping.getSqlType()) { case Types.BOOLEAN: return resultSet.getBoolean(mapping.getColumnIndex()); case Types.TINYINT: case Types.SMALLINT: case Types.INTEGER: case Types.BIGINT: case Types.ROWID: return resultSet.getInt(mapping.getColumnIndex()); case Types.CHAR: case Types.VARCHAR: case Types.LONGVARCHAR: case Types.NCHAR: case Types.NVARCHAR: case Types.LONGNVARCHAR: case Types.SQLXML: return resultSet.getString(mapping.getColumnIndex()); case Types.REAL: case Types.FLOAT: return resultSet.getFloat(mapping.getColumnIndex()); case Types.DOUBLE: return resultSet.getDouble(mapping.getColumnIndex()); case Types.NUMERIC: return resultSet.getBigDecimal(mapping.getColumnIndex()); case Types.DECIMAL: return resultSet.getBigDecimal(mapping.getColumnIndex()); case Types.DATE: return resultSet.getDate(mapping.getColumnIndex()); case Types.TIME: case Types.TIME_WITH_TIMEZONE: return resultSet.getTime(mapping.getColumnIndex()); case Types.TIMESTAMP: case Types.TIMESTAMP_WITH_TIMEZONE: return resultSet.getTimestamp(mapping.getColumnIndex()); case Types.BINARY: case Types.VARBINARY: case Types.LONGVARBINARY: case Types.NULL: case Types.OTHER: case Types.JAVA_OBJECT: case Types.DISTINCT: case Types.STRUCT: case Types.ARRAY: case Types.BLOB: case Types.CLOB: case Types.REF: case Types.DATALINK: case Types.NCLOB: case Types.REF_CURSOR: return resultSet.getByte(mapping.getColumnIndex()); default: return resultSet.getString(mapping.getColumnIndex()); }
ResultSetSchemaGenerator类(原始版本)
public class ResultSetSchemaGenerator { /** * Generates parquet schema using {@link ResultSetMetaData } * * @param resultSet * @param name * Record name. * @param nameSpace * @return * @throws SQLException */ public SchemaResults generateSchema(ResultSet resultSet, String name, String nameSpace) throws SQLException { SchemaResults schemaResults = new SchemaResults(); List<SchemaSqlMapping> mappings = new ArrayList<>(); Schema recordSchema = Schema.createRecord(name, null, nameSpace, false); List<Schema.Field> fields = new ArrayList<>(); if (resultSet != null) { ResultSetMetaData resultSetMetaData = resultSet.getMetaData(); int columnCount = resultSetMetaData.getColumnCount(); for (int x = 1; x <= columnCount; x++) { String columnName = resultSetMetaData.getColumnName(x).replaceAll("[^a-zA-Z0-9_]", ""); int sqlColumnType = resultSetMetaData.getColumnType(x); String schemaName = columnName.toLowerCase(); Schema.Type schemaType = parseSchemaType(sqlColumnType); mappings.add(new SchemaSqlMapping(schemaName, columnName, sqlColumnType, x, schemaType)); fields.add(createNullableField(schemaName, schemaType)); } } recordSchema.setFields(fields); schemaResults.setMappings(mappings); schemaResults.setParsedSchema(recordSchema); return schemaResults; } public Schema.Type parseSchemaType(int sqlColumnType) { switch (sqlColumnType) { case Types.BOOLEAN: return Schema.Type.BOOLEAN; case Types.TINYINT: // 1 byte case Types.SMALLINT: // 2 bytes case Types.INTEGER: // 4 bytes return Schema.Type.INT; // 32 bit (4 bytes) (signed) case Types.ROWID: case Types.CHAR: case Types.VARCHAR: case Types.LONGVARCHAR: case Types.NCHAR: case Types.NVARCHAR: case Types.LONGNVARCHAR: case Types.SQLXML: return Schema.Type.STRING; // unicode string case Types.REAL: // Approximate numerical (mantissa single precision 7) return Schema.Type.FLOAT; // A 32-bit IEEE single-float case Types.DOUBLE: // Approximate numerical (mantissa precision 16) case Types.DECIMAL: // Exact numerical (5 - 17 bytes) case Types.NUMERIC: // Exact numerical (5 - 17 bytes) case Types.FLOAT: // Approximate numerical (mantissa precision 16) return Schema.Type.DOUBLE; // A 64-bit IEEE double-float case Types.DATE: case Types.TIME: case Types.TIMESTAMP: case Types.TIME_WITH_TIMEZONE: case Types.TIMESTAMP_WITH_TIMEZONE: case Types.BIGINT: // 8 bytes return Schema.Type.LONG; // 64 bit (signed) case Types.BINARY: case Types.VARBINARY: case Types.LONGVARBINARY: case Types.NULL: case Types.OTHER: case Types.JAVA_OBJECT: case Types.DISTINCT: case Types.STRUCT: case Types.ARRAY: case Types.BLOB: case Types.CLOB: case Types.REF: case Types.DATALINK: case Types.NCLOB: case Types.REF_CURSOR: return Schema.Type.BYTES; // sequence of bytes } return null; } /** * Creates field that can be null {@link Schema#createUnion(List)} * * @param columnName * The schema column name * @param type * Schema type for this field * @return */ public Schema.Field createNullableField(String columnName, Schema.Type type) { Schema intSchema = Schema.create(type); Schema nullSchema = Schema.create(Schema.Type.NULL); List<Schema> fieldSchemas = new ArrayList<>(); fieldSchemas.add(intSchema); fieldSchemas.add(nullSchema); Schema fieldSchema = Schema.createUnion(fieldSchemas); return new Schema.Field(columnName, fieldSchema, null, null); } }
问题根源
当前的ResultSetSchemaGenerator仅将SQL的Timestamp/Date/Time类型映射为Avro的LONG类型,但未添加对应的LogicalType元数据。Avro的LogicalTypeConversion是基于Schema中的logicalType属性触发的,缺少该属性时转换逻辑不会生效,导致代码尝试直接将java.sql.Timestamp强制转为Long,从而抛出类型转换异常。
具体修复步骤
1. 修改Schema生成逻辑,添加LogicalType属性
修改ResultSetSchemaGenerator的createNullableField方法,接收原始SQL列类型参数,根据不同的SQL时间类型给Avro Schema添加对应的logicalType属性:
- Types.DATE → "date"
- Types.TIME → "time-millis"
- Types.TIME_WITH_TIMEZONE → "time-micros"
- Types.TIMESTAMP → "timestamp-millis"
- Types.TIMESTAMP_WITH_TIMEZONE → "timestamp-micros"
同时调整generateSchema方法中调用createNullableField的代码,传入sqlColumnType参数。
修改后的ResultSetSchemaGenerator关键代码:
// 修改generateSchema中的调用 fields.add(createNullableField(schemaName, schemaType, sqlColumnType)); // 修改createNullableField方法 public Schema.Field createNullableField(String columnName, Schema.Type type, int sqlColumnType) { Schema baseSchema = Schema.create(type); // 添加LogicalType元数据 switch (sqlColumnType) { case Types.DATE: baseSchema.addProp("logicalType", "date"); break; case Types.TIME: baseSchema.addProp("logicalType", "time-millis"); break; case Types.TIME_WITH_TIMEZONE: baseSchema.addProp("logicalType", "time-micros"); break; case Types.TIMESTAMP: baseSchema.addProp("logicalType", "timestamp-millis"); break; case Types.TIMESTAMP_WITH_TIMEZONE: baseSchema.addProp("logicalType", "timestamp-micros"); break; default: // 其他类型无需处理 } Schema nullSchema = Schema.create(Schema.Type.NULL); List<Schema> fieldSchemas = new ArrayList<>(); fieldSchemas.add(baseSchema); fieldSchemas.add(nullSchema); Schema fieldSchema = Schema.createUnion(fieldSchemas); return new Schema.Field(columnName, fieldSchema, null, null); }
2. 确认DataModel配置有效性
当前添加的TimeConversions已覆盖常用时间类型转换,无需修改,但要确保AvroParquetWriter正确使用了配置好的timeSupport DataModel。
3. 保留结果集提取逻辑
extractResult方法中返回的java.sql.Timestamp/Date/Time对象,会被对应的LogicalTypeConversion自动转换为Avro要求的数值类型(例如timestamp-millis对应毫秒级时间戳Long值),无需手动转换。
验证效果
修改后,Avro Schema中的时间字段会带有logicalType属性,例如Timestamp列的Schema示例:
{ "name": "create_time", "type": ["long", "null"], "logicalType": "timestamp-millis" }
此时GenericData的转换逻辑会被触发,自动完成java.sql.Timestamp到Long时间戳的转换,避免类型转换异常。
内容的提问来源于stack exchange,提问作者javaseeker

