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

Scala按字节拆分字符串发Kinesis 代码运行抛出空指针异常

问题根因排查

你现在的代码有两个核心问题,其中第一个直接触发空指针异常:

  • 空指针触发点:foreach循环内的finally块会在第一个分片发送完成后立刻执行kinesis.close(),Kinesis客户端被销毁后,后续分片调用kinesis.putRecord时会直接访问已释放的资源,抛出空指针。
  • 逻辑错误点:String.grouped(1048576)是按字符个数拆分,不是按字节数拆分。对于中文、emoji等多字节UTF-8字符,单块实际字节长度会远超1MB的Kinesis单条记录上限,既不符合需求,也会触发Kinesis的参数校验错误。
  • 额外风险点:代码没有做前置空值校验,如果dataArray、kinesis实例、streamName任意一个为null,也会抛出空指针。
正确实现方案

修正逻辑遵循三个原则:客户端资源在全部分片发送完成后统一释放、按字节拆分时不截断多字节字符避免乱码、前置空校验拦截非法输入。

import java.nio.charset.StandardCharsets
import software.amazon.awssdk.core.SdkBytes
import scala.collection.mutable.ArrayBuffer

// 前置空值校验
require(dataArray != null, "待发送数据数组不能为null")
require(kinesis != null, "Kinesis客户端实例不能为null")
require(streamName != null && streamName.nonEmpty, "流名称不能为null或空")

val recordDelimiter: String = "\n"
val recordsBatch: String = dataArray.mkString(recordDelimiter)
val recordsCount = dataArray.length
val maxChunkBytes = 1048576
val charset = StandardCharsets.UTF_8

/**
 * 按指定字节上限拆分字符串,保证不截断多字节字符
 */
def splitByByteSize(source: String, maxBytes: Int, cs: java.nio.charset.Charset): List[String] = {
  val encoder = cs.newEncoder()
  val sourceBuffer = cs.encode(source)
  val chunks = ArrayBuffer.empty[String]

  while (sourceBuffer.hasRemaining) {
    val currentChunkSize = Math.min(maxBytes, sourceBuffer.remaining())
    val sliceBuffer = sourceBuffer.slice()
    sliceBuffer.limit(currentChunkSize)

    // 回退到合法字符边界,避免截断多字节字符导致乱码
    var validEndPos = currentChunkSize
    while (validEndPos > 0 && !encoder.isLegalReplacement(sliceBuffer.slice(0, validEndPos))) {
      validEndPos -= 1
    }

    val chunkBytes = new Array[Byte](validEndPos)
    sliceBuffer.get(chunkBytes, 0, validEndPos)
    chunks += new String(chunkBytes, cs)
    sourceBuffer.position(sourceBuffer.position() + validEndPos)
  }
  chunks.toList
}

// 全部分片发送完成后再释放客户端资源
try {
  splitByByteSize(recordsBatch, maxChunkBytes, charset)
    .foreach { chunk =>
      try {
        val putRecordRequest = PutRecordRequest.builder()
          .streamName(streamName)
          .data(SdkBytes.fromString(chunk, charset))
          .partitionKey("partitionKey-1")
          .build()
        val response = kinesis.putRecord(putRecordRequest)
        logger.info(s"Records sent to $streamName stream, total records count: $recordsCount, sequence number: ${response.sequenceNumber()}")
      } catch {
        case exception: KinesisException =>
          logger.error("Error occurred while sending data to Kinesis", exception)
      }
    }
} finally {
  kinesis.close()
}

优化建议

  • 如果需要保证单条记录不被拆分到不同分片,可以不用先拼接全量字符串,直接遍历dataArray累加单条记录的字节长度,凑到接近1MB阈值时就打包发送,性能比全量拼接后拆分更高。
  • 高吞吐量场景下建议用Kinesis的putRecords批量接口替代单条putRecord,可以大幅降低请求开销。

内容的提问来源于stack exchange,提问作者MasterC

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 21:36:25