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
相关产品推荐
相关产品推荐

