如何在不修改字节的前提下将字符串转浮点数?修复Kafka数据编码问题
解决Debezium Decimal转S3后字符串转数值的Python 2.7方案
我之前碰到过一模一样的问题,Debezium把MySQL的Decimal类型序列化成Avro的decimal逻辑类型后,写入S3时字节被当成ASCII字符串保存,读出来就是JiU8这类看起来像乱码的字符串,其实只要逆向解析就能拿到正确的数值。
问题本质
Debezium的Decimal序列化规则是把数值转换成大端字节序的带符号整数补码(对应你提供的Avro Schema里的bytes类型+decimal逻辑类型),而写入S3时这些字节被直接当作ASCII字符存储成了字符串。我们要做的就是把这个字符串还原成字节,解析成整数,再除以10的scale次方(你的schema里scale是2,对应两位小数)。
转换代码(Python 2.7)
根据你Schema里precision=64的设置,我们用64位大端整数来解析:
import struct def decode_debezium_decimal(decimal_str, scale=2): # Python 2.7中字符串本身就是字节序列,直接使用即可 byte_data = decimal_str # 解析为大端模式的64位带符号整数 raw_integer = struct.unpack('>q', byte_data)[0] # 转换为带小数位的数值 return float(raw_integer) / (10 ** scale)
验证示例
把你给出的测试值代入,完全符合预期:
print(decode_debezium_decimal('JiU8')) # 输出 24999.0 print(decode_debezium_decimal('JiDw')) # 输出 24988.0 print(decode_debezium_decimal('RxFc')) # 输出 46575.0 print(decode_debezium_decimal('LyZQ')) # 输出 30900.0
为什么你之前的尝试不对?
你用struct.unpack('f', ...)是把字节当作单精度浮点数的二进制格式解析,但Debezium存储的是整数的二进制补码,不是浮点数的二进制表示,所以肯定得不到正确结果。必须先解析成整数再转小数才对。
在Spark任务中批量转换
如果你要在Python 2.7的Spark任务里批量处理这个字段,可以注册一个UDF:
from pyspark.sql.functions import udf from pyspark.sql.types import FloatType # 注册转换UDF decimal_decode_udf = udf(decode_debezium_decimal, FloatType()) # 处理DataFrame中的PRICE_SELLING字段 processed_df = original_df.withColumn( "PRICE_SELLING", decimal_decode_udf(original_df.PRICE_SELLING) )
内容的提问来源于stack exchange,提问作者ElMoselYEE
相关产品推荐
相关产品推荐

