Spring Boot应用无法读取Kinesis流,调用getRecords时JSON解析异常
Kinesis流读取失败及JSON解析异常排查方案
核心问题分析
出现SdkClientException且JSON解析时遇到意外字符¿(编码191),大概率是Kinesis响应的编码格式不匹配,或是AWS客户端配置的序列化/反序列化器存在问题;同时DynamoDB无条目说明KCL的checkpoint机制未正常初始化或运行。
具体排查步骤
1. 检查Kinesis客户端的序列化配置
KCL 3.x默认使用Jackson处理JSON序列化,需确保客户端配置中指定了正确的字符编码:
// 示例:配置Kinesis客户端时显式指定UTF-8编码 SdkHttpClient httpClient = ApacheHttpClient.builder() .connectionTimeout(Duration.ofSeconds(10)) .socketTimeout(Duration.ofSeconds(10)) .build(); KinesisAsyncClient kinesisClient = KinesisAsyncClient.builder() .httpClient(httpClient) .region(Region.US_EAST_1) .overrideConfiguration(cfg -> cfg.addExecutionInterceptor(new ExecutionInterceptor() { @Override public void beforeExecution(Context.BeforeExecution context) { context.request().headers().put("Content-Type", "application/json; charset=utf-8"); } })) .build();
2. 验证Kinesis流中数据的编码格式
流中的记录可能非UTF-8编码,导致反序列化失败。可以通过AWS CLI先手动读取一条记录验证:
aws kinesis get-records --shard-iterator <你的shard迭代器> --limit 1
查看返回的Data字段(Base64编码),解码后检查是否包含非UTF-8字符,或是否为预期的JSON格式。
3. 检查KCL的Checkpoint配置
DynamoDB无条目说明KCL未成功创建checkpoint表,需确认:
- AWS账户对DynamoDB有
CreateTable、PutItem等权限 - KCL配置中指定了正确的应用名称(
applicationName),该名称对应DynamoDB表名 - 客户端配置的Region与Kinesis流、DynamoDB的Region一致
4. 自定义记录处理器的反序列化逻辑
如果流中数据不是标准JSON,需在记录处理器中手动处理解码:
public class CustomRecordProcessor implements RecordProcessor { @Override public void processRecords(ProcessRecordsInput processRecordsInput) { for (Record record : processRecordsInput.records()) { try { // 先解码Base64,再指定UTF-8处理字符 String data = new String(Base64.getDecoder().decode(record.data()), StandardCharsets.UTF_8); // 后续处理逻辑 } catch (Exception e) { // 记录错误日志,跳过异常记录避免阻塞 LOGGER.error("处理记录失败", e); } } } }
5. 排查依赖冲突
Java 11下需确保Jackson版本与KCL 3.0.3兼容,避免依赖冲突。可通过./gradlew dependencies(Gradle)或mvn dependency:tree(Maven)检查Jackson相关依赖的版本,统一使用KCL依赖的Jackson版本。
内容的提问来源于stack exchange,提问作者Rohan Sharma
相关产品推荐
相关产品推荐

