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

使用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返回的空列表时,就会触发这个异常。

修复步骤

  1. 替换不可序列化的列表实现:在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());
        }
}
  1. 优化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 03:35:26