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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 07:00:54