You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 14:30:59