如何让ThriftParquetWriter将16位int以short写入而非转32位int
我太懂这种需求了——明明在Thrift类里定义的是i16字段,用ThriftParquetWriter写入Parquet后,却被自动转成32位整数存储,完全不符合预期。虽然官方文档提到32位整数的高效编码能覆盖16位的场景,但咱就是要精准存储16位整数,对吧?下面给你一步步拆解解决方案,包括最便捷的注解方式,以及兜底的自定义转换器方案。
一、最推荐:用Thrift注解指定Parquet逻辑类型
ThriftParquetWriter支持通过Thrift IDL的字段注解,直接映射到Parquet的逻辑类型。我们只需要给i16字段加上专属注解,明确告诉序列化器要把它存成INT16类型。
修改你的Thrift IDL文件
在你的Thrift结构体定义中,给目标i16字段添加@parquet.type("INT16")注解:
namespace java com.mycompany struct MyClass { // 给i16字段标记Parquet逻辑类型为INT16 @parquet.type("INT16") 1: optional i16 shortField; // 其他字段正常定义... }
这个注解是核心,它会覆盖默认的类型映射规则,强制Parquet把这个i16字段序列化为16位整数类型。
确保依赖版本支持
要注意,这个注解功能需要parquet-thrift 1.10及以上版本(建议用1.12+的稳定版),所以检查你的构建依赖:
比如Maven依赖:
<dependency> <groupId>org.apache.parquet</groupId> <artifactId>parquet-thrift</artifactId> <version>1.14.0</version> </dependency>
二、兜底方案:自定义ThriftRecordConverter
如果注解方式在你的环境中不生效(比如用了较旧的库版本),可以通过自定义转换器来手动控制i16字段的序列化逻辑。
1. 实现自定义转换器
创建一个继承ThriftRecordConverter的类,重写字段写入逻辑,专门处理i16类型:
import org.apache.parquet.thrift.ThriftRecordConverter; import org.apache.parquet.thrift.ThriftSchemaConverter; import org.apache.parquet.schema.MessageType; import org.apache.parquet.thrift.struct.ThriftField; import org.apache.parquet.thrift.struct.ThriftFieldType; public class CustomShortFieldConverter extends ThriftRecordConverter { public CustomShortFieldConverter(MessageType parquetSchema, ThriftSchemaConverter schemaConverter) { super(parquetSchema, schemaConverter); } @Override protected void writeField(ThriftField field, Object value) { // 只针对i16字段做特殊处理 if (field.getType() == ThriftFieldType.I16) { // 直接写入16位整数,跳过默认的32位转换 getRecordConsumer().addInteger((Short) value); } else { // 其他字段沿用默认序列化逻辑 super.writeField(field, value); } } }
2. 自定义ThriftParquetWriter
创建自定义的Writer类,指定使用我们刚才的转换器:
import org.apache.hadoop.conf.Configuration; import org.apache.parquet.hadoop.ParquetWriter; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import org.apache.parquet.thrift.ThriftParquetWriter; import org.apache.parquet.thrift.ThriftSchemaConverter; import org.apache.parquet.schema.MessageType; import org.apache.thrift.TBase; import java.io.IOException; public class CustomThriftParquetWriter<T extends TBase> extends ThriftParquetWriter<T> { public CustomThriftParquetWriter(String path, Class<T> thriftClass, CompressionCodecName compressionCodecName, int blockSize, int pageSize, boolean enableDictionary, boolean validating, Configuration conf) throws IOException { super(path, thriftClass, compressionCodecName, blockSize, pageSize, enableDictionary, validating, conf); } @Override protected ThriftRecordConverter createThriftRecordConverter(MessageType parquetSchema, ThriftSchemaConverter schemaConverter) { // 返回自定义的转换器 return new CustomShortFieldConverter(parquetSchema, schemaConverter); } }
3. 使用自定义Writer写入数据
替换原来的ThriftParquetWriter为我们的自定义类即可:
CustomThriftParquetWriter writer = new CustomThriftParquetWriter( outPutPathInS3, Class.forName("com.mycompany.MyClass"), CompressionCodecName.SNAPPY, BLOCK_SIZE, PAGE_SIZE, ParquetWriter.DEFAULT_IS_DICTIONARY_ENABLED, ParquetWriter.DEFAULT_IS_VALIDATING_ENABLED, conf ); for (TBase trecord : recordsToFlush) { writer.write(trecord); }
三、验证结果
写入完成后,可以用parquet-tools工具查看Parquet文件的Schema,确认目标字段的类型是INT16:
parquet-tools schema your-target-file.parquet
如果看到类似optional int32 shortField (INTEGER(16,true));或者直接标注INT16的字段,就说明转换成功了。
内容的提问来源于stack exchange,提问作者Viacheslav Shalamov

