Scala中压缩java.util.List[java.lang.Float]的方案咨询
Scala环境下压缩java.util.List[java.lang.Float]的实用方案
针对你的Flink流处理场景(数百个java.util.List[java.lang.Float]需压缩以降低内存占用),以下是几个可落地的方案:
1. 原生类型转换+通用压缩算法
先将包装类型的List[java.lang.Float]转为原生float[],直接削减对象头带来的额外内存开销,再用通用压缩算法压缩字节数组。这是最基础且普适的方案。
代码示例
import java.util.zip.{Deflater, Inflater} import java.util.{List => JList} object FloatListCompressor { // 压缩方法 def compress(list: JList[java.lang.Float]): Array[Byte] = { if (list.isEmpty) return Array.emptyByteArray // 转原生float数组 val floatArr = new Array[Float](list.size()) var idx = 0 while (idx < list.size()) { floatArr(idx) = list.get(idx) idx += 1 } // 转字节数组 val byteBuffer = java.nio.ByteBuffer.allocate(floatArr.length * 4) byteBuffer.asFloatBuffer().put(floatArr) val rawBytes = byteBuffer.array() // Deflater压缩 val deflater = new Deflater() deflater.setInput(rawBytes) deflater.finish() val compressedBuf = new Array[Byte](rawBytes.length) val compressedLen = deflater.deflate(compressedBuf) deflater.reset() // 截取有效压缩结果 java.util.Arrays.copyOf(compressedBuf, compressedLen) } // 解压方法 def decompress(compressed: Array[Byte]): JList[java.lang.Float] = { val inflater = new Inflater() inflater.setInput(compressed) val decompressedBuf = new Array[Byte](compressed.length * 4) // 预估解压长度 val decompressedLen = inflater.inflate(decompressedBuf) inflater.end() val byteBuffer = java.nio.ByteBuffer.wrap(decompressedBuf, 0, decompressedLen) val floatBuffer = byteBuffer.asFloatBuffer() val floatArr = new Array[Float](floatBuffer.limit()) floatBuffer.get(floatArr) // 转回java.util.List java.util.Arrays.asList(floatArr.map(java.lang.Float.valueOf): _*) } }
2. 差值编码+压缩(适合连续变化的数值)
如果你的Float列表是连续采样数据(如传感器读数),先对列表做差值编码:存储第一个元素的原始值,后续元素存储与前一个元素的差值。差值的数值范围通常远小于原始值,可进一步转成更小的原生类型(如short)再压缩,能获得更高的压缩率。
代码示例
def compressWithDelta(list: JList[java.lang.Float]): Array[Byte] = { if (list.isEmpty) return Array.emptyByteArray val floatArr = new Array[Float](list.size()) var idx = 0 while (idx < list.size()) { floatArr(idx) = list.get(idx) idx += 1 } // 差值编码 val deltaArr = new Array[Float](floatArr.length) deltaArr(0) = floatArr(0) idx = 1 while (idx < floatArr.length) { deltaArr(idx) = floatArr(idx) - floatArr(idx - 1) idx += 1 } // 将差值转为short(需确保差值在Short范围内) val shortArr = deltaArr.tail.map(_.toShort) // 组装字节数组:首元素(4字节)+ 差值数组(每个2字节) val byteBuffer = java.nio.ByteBuffer.allocate(4 + shortArr.length * 2) byteBuffer.putFloat(deltaArr(0)) byteBuffer.asShortBuffer().put(shortArr) val rawBytes = byteBuffer.array() // 后续压缩步骤同方案1 val deflater = new Deflater() deflater.setInput(rawBytes) deflater.finish() val compressedBuf = new Array[Byte](rawBytes.length) val compressedLen = deflater.deflate(compressedBuf) deflater.reset() java.util.Arrays.copyOf(compressedBuf, compressedLen) }
3. 高速压缩库(适合流处理低延迟需求)
流处理场景对CPU开销敏感,可选用LZ4、Snappy这类高速压缩库,它们的压缩/解压速度远高于Deflater,同时能提供不错的压缩率。以LZ4为例:
代码示例(需引入LZ4依赖)
import net.jpountz.lz4.LZ4Factory import java.util.{List => JList} def compressWithLZ4(list: JList[java.lang.Float]): Array[Byte] = { val floatArr = new Array[Float](list.size()) var idx = 0 while (idx < list.size()) { floatArr(idx) = list.get(idx) idx += 1 } val byteBuffer = java.nio.ByteBuffer.allocate(floatArr.length * 4) byteBuffer.asFloatBuffer().put(floatArr) val rawBytes = byteBuffer.array() // LZ4高速压缩 val lz4 = LZ4Factory.fastestInstance().fastCompressor() val maxCompressedLen = lz4.maxCompressedLength(rawBytes.length) val compressedBuf = new Array[Byte](maxCompressedLen) val compressedLen = lz4.compress(rawBytes, 0, rawBytes.length, compressedBuf, 0, maxCompressedLen) java.util.Arrays.copyOf(compressedBuf, compressedLen) }
4. Flink状态层优化
如果这些列表是作为Flink状态存储的,可结合Flink自身的状态优化能力:
- 优先用原生
float[]而非java.util.List[java.lang.Float]存储状态,削减包装类型的内存开销; - 配置RocksDB状态后端启用压缩(如Snappy/LZ4),全局压缩所有状态数据,无需手动处理压缩逻辑;
- 若使用
ListState,可考虑将多个小列表合并为大数组后再存储,减少状态元数据开销。
状态配置示例
在Flink作业配置中开启RocksDB压缩:
import org.apache.flink.contrib.streaming.state.RocksDBStateBackend import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment val env = StreamExecutionEnvironment.getExecutionEnvironment() val rocksDBBackend = new RocksDBStateBackend("hdfs:///path/to/rocksdb") // 开启Snappy压缩 rocksDBBackend.setCompressionType(org.rocksdb.CompressionType.SNAPPY_COMPRESSION) env.setStateBackend(rocksDBBackend)
内容的提问来源于stack exchange,提问作者Alok Singh
相关产品推荐
相关产品推荐

