使用Apache Beam DynamoDBIO读取DynamoDB指定数据时序列化异常求助
问题
我用Apache Beam的DynamoDBIO SDK实现了一个从DynamoDB读取数据的管道,需要通过filterExpression过滤指定数据,当前代码如下:
Map<String, AttributeValue> expressionAttributeValues = new HashMap<>(); expressionAttributeValues.put(":message", AttributeValue.builder().s("Ping").build()); pipeline .apply(DynamoDBIO.<List<Map<String, AttributeValue>>>read() .withClientConfiguration(DynamoDBConfig.CLIENT_CONFIGURATION) .withScanRequestFn(input -> ScanRequest.builder().tableName("SiteProductCache").totalSegments(1) .filterExpression("KafkaEventMessage = :message") .expressionAttributeValues(expressionAttributeValues) .projectionExpression("key, KafkaEventMessage") .build()) .withScanResponseMapperFn(new ResponseMapper()) .withCoder(ListCoder.of(MapCoder.of(StringUtf8Coder.of(), AttributeValueCoder.of()))) ) .apply(...) ---- static final class ResponseMapper implements SerializableFunction<ScanResponse, List<Map<String, AttributeValue>>> { @Override public List<Map<String, AttributeValue>> apply(ScanResponse input) { if (input == null) { return Collections.emptyList(); } return input.items(); } }
执行时抛出以下异常:
Exception in thread "main" java.lang.IllegalArgumentException: Forbidden IOException when writing to OutputStream at org.apache.beam.sdk.util.CoderUtils.encodeToSafeStream(CoderUtils.java:89) at org.apache.beam.sdk.util.CoderUtils.encodeToByteArray(CoderUtils.java:70) at org.apache.beam.sdk.util.CoderUtils.encodeToByteArray(CoderUtils.java:55) at org.apache.beam.sdk.transforms.Create$Values$CreateSource.fromIterable(Create.java:413) at org.apache.beam.sdk.transforms.Create$Values.expand(Create.java:370) at org.apache.beam.sdk.transforms.Create$Values.expand(Create.java:277) at org.apache.beam.sdk.Pipeline.applyInternal(Pipeline.java:548) at org.apache.beam.sdk.Pipeline.applyTransform(Pipeline.java:499) at org.apache.beam.sdk.values.PBegin.apply(PBegin.java:56) at org.apache.beam.sdk.io.aws2.dynamodb.DynamoDBIO$Read.expand(DynamoDBIO.java:301) at org.apache.beam.sdk.io.aws2.dynamodb.DynamoDBIO$Read.expand(DynamoDBIO.java:172) at org.apache.beam.sdk.Pipeline.applyInternal(Pipeline.java:548) at org.apache.beam.sdk.Pipeline.applyTransform(Pipeline.java:482) at org.apache.beam.sdk.values.PBegin.apply(PBegin.java:44) at org.apache.beam.sdk.Pipeline.apply(Pipeline.java:177) at some_package.beam_state_storage.dynamodb.DynamoDBPipelineDefinition.run(DynamoDBPipelineDefinition.java:40) at some_package.beam_state_storage.dynamodb.DynamoDBPipelineDefinition.main(DynamoDBPipelineDefinition.java:28) Caused by: java.io.NotSerializableException: software.amazon.awssdk.core.util.DefaultSdkAutoConstructList at java.base/java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1197) at java.base/java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1582) at java.base/java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1539) at java.base/java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1448) Caused by: java.io.NotSerializableException: software.amazon.awssdk.core.util.DefaultSdkAutoConstructList
解决方案
异常原因
DefaultSdkAutoConstructList是AWS SDK内部的空列表实现,不支持Java序列化。而Beam在分布式执行时需要序列化所有传递的对象,当ResponseMapper直接返回AWS返回的空列表时,就会触发这个异常。
修复步骤
- 替换不可序列化的列表实现:在
ResponseMapper中,将返回的列表转换为可序列化的ArrayList,即使是空结果也手动创建新的ArrayList实例:
static final class ResponseMapper implements SerializableFunction<ScanResponse, List<Map<String, AttributeValue>>> { @Override public List<Map<String, AttributeValue>> apply(ScanResponse input) { if (input == null || input.items() == null) { return new ArrayList<>(); } // 将AWS返回的列表包装为可序列化的ArrayList return new ArrayList<>(input.items()); } }
- 优化ScanRequest的创建逻辑:把
expressionAttributeValues的创建移到withScanRequestFn的lambda内部,避免外部对象的引用可能带来的序列化问题(虽然当前Map是可序列化的,但这种写法更符合Beam的分布式执行规范):
pipeline .apply(DynamoDBIO.<List<Map<String, AttributeValue>>>read() .withClientConfiguration(DynamoDBConfig.CLIENT_CONFIGURATION) .withScanRequestFn(input -> { Map<String, AttributeValue> expressionAttributeValues = new HashMap<>(); expressionAttributeValues.put(":message", AttributeValue.builder().s("Ping").build()); return ScanRequest.builder() .tableName("SiteProductCache") .totalSegments(1) .filterExpression("KafkaEventMessage = :message") .expressionAttributeValues(expressionAttributeValues) .projectionExpression("key, KafkaEventMessage") .build(); }) .withScanResponseMapperFn(new ResponseMapper()) .withCoder(ListCoder.of(MapCoder.of(StringUtf8Coder.of(), AttributeValueCoder.of()))) ) .apply(...)
额外建议
- 如果过滤条件基于主键或GSI(全局二级索引),优先使用
Query操作替代Scan,Query的性能远优于全表扫描,Beam的DynamoDBIO同样支持withQueryRequestFn方法。 - 确保所有在Beam转换中使用的自定义类、捕获的外部对象都实现
Serializable接口,避免类似的序列化异常。
内容的提问来源于stack exchange,提问作者Prasad
相关产品推荐
相关产品推荐

