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

如何用Moq在.NET中Mock AWS S3的SelectObjectContentAsync方法

问题:Mock S3的SelectObjectContentAsync方法时出现校验和错误

尝试Mock S3的SelectObjectContentAsync方法,找到了Java和Go的示例,但没有C#的。之前成功Mock过GetObjectAsync,但这次遇到了校验和错误,推测是Mock配置问题。


错误信息:
System.Net.Http.HttpRequestException : InternalServerError - Internal Server Error: Amazon.S3.Model.S3EventStreamException: Error.
---> Amazon.Runtime.EventStreams.EventStreamChecksumFailureException: Message Prelude Checksum failure. Expected 1382380405 but was -1954053792
at Amazon.Runtime.EventStreams.Internal.EventStreamDecoder.ProcessPrelude(Byte[] data, Int32 offset, Int32 length)
at Amazon.Runtime.EventStreams.Internal.EventStreamDecoder.ProcessData(Byte[] data, Int32 offset, Int32 length)
at Amazon.Runtime.EventStreams.Internal.EventStream2.ReadFromStream(Byte[] buffer) at Amazon.Runtime.EventStreams.Internal.EnumerableEventStream2.GetEnumerator()+MoveNext()
--- End of inner exception stack trace ---
at Amazon.Runtime.EventStreams.Internal.EnumerableEventStream`2.GetEnumerator()+MoveNext()


原Mock S3客户端代码

using Amazon.S3;
using Amazon.S3.Model;
using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Moq;
class MockIAmazonS3
{
    private static Mock<IAmazonS3> GetMockS3()
    {
        var mockS3Client = new Mock<IAmazonS3>();            
        mockS3Client.Setup(x =>
            x.SelectObjectContentAsync(It.IsAny<SelectObjectContentRequest>(), It.IsAny<CancellationToken>()))
        .Returns(new Func<SelectObjectContentRequest, CancellationToken, Task<SelectObjectContentResponse>>(GetSelectObjectContentResponse));

        return mockS3Client;
    }
    private static async Task<SelectObjectContentResponse> GetSelectObjectContentResponse(SelectObjectContentRequest request,
    CancellationToken cancellationToken = default)
    {
        var streamObj = new MemoryStream();
        // 根据request.Key设置json字符串
        var jsonString = "{a:1,b:1},{a:2,b:2}";
        await System.Text.Json.JsonSerializer.SerializeAsync(streamObj, jsonString);
        streamObj.Position = 0;
        var response = new SelectObjectContentResponse {
            Payload = new SelectObjectContentEventStream(streamObj ) {
            }
        };
        return response;
    }
}

业务逻辑代码

private async Task QueryS3Select<T>(string bucket, string fileKey, string expression) where T : new()
{
    var s3client = await s3.CreateClient(bucket);
    SelectObjectContentRequest request = new SelectObjectContentRequest();
    request.BucketName = bucket;
    request.Key = fileKey;
    request.Expression = $"SELECT * FROM s3object {expression}";
    request.ExpressionType = ExpressionType.SQL;
    request.InputSerialization = new InputSerialization
    {
        Parquet = new ParquetInput()
    };
    request.OutputSerialization = new OutputSerialization
    {
        JSON = new JSONOutput
        {
            RecordDelimiter = ","
        }
    };
    SelectObjectContentResponse response = await s3client.SelectObjectContentAsync(request);
    using (var eventStream = response.Payload)
    {
        eventStream.ExceptionReceived += (sender, args) => throw args.EventStreamException;
        var recordResults = eventStream
            .Where(ev => ev is RecordsEvent)
            .Cast<RecordsEvent>()
            .Select(records =>
            {
                using (var reader = new StreamReader(records.Payload, Encoding.UTF8))
                {
                    return reader.ReadToEnd();
                }
            }).ToList();
    }
}

解决方案:自定义Mock事件流实现

需要Mock ISelectObjectContentEventStream 接口,自定义实现类:

internal class MockSelectObjectContentEventStream : ISelectObjectContentEventStream
{
    private readonly List<IS3Event> s3Events;
    public int BufferSize { get => throw new NotImplementedException(); set => throw new NotImplementedException(); }

    public event EventHandler<EventStreamEventReceivedArgs<IS3Event>> EventReceived;
    public event EventHandler<EventStreamExceptionReceivedArgs<S3EventStreamException>> ExceptionReceived;
    public event EventHandler<EventStreamEventReceivedArgs<RecordsEvent>> RecordsEventReceived;
    public event EventHandler<EventStreamEventReceivedArgs<StatsEvent>> StatsEventReceived;
    public event EventHandler<EventStreamEventReceivedArgs<ProgressEvent>> ProgressEventReceived;
    public event EventHandler<EventStreamEventReceivedArgs<ContinuationEvent>> ContinuationEventReceived;
    public event EventHandler<EventStreamEventReceivedArgs<EndEvent>> EndEventReceived;

    public MockSelectObjectContentEventStream(List<IS3Event> S3Events)
    {
        s3Events = S3Events;
    }
    public void Dispose()
    {
    }
    public IEnumerator<IS3Event> GetEnumerator()
    {
        return s3Events.GetEnumerator();
    }

    public void StartProcessing()
    {
    }

    public Task StartProcessingAsync()
    {
        return Task.CompletedTask;
    }

    IEnumerator IEnumerable.GetEnumerator()
    {
        return s3Events.GetEnumerator();
    }
}

修改Mock方法,返回自定义事件流的响应:

private static async Task<SelectObjectContentResponse> GetSelectObjectContentResponse(SelectObjectContentRequest request,
    CancellationToken cancellationToken = default)
{
    var jsonstring = "{}";
    switch (request.Key) {
        case "request1":
            jsonstring = "{a:1}";
            break;
        case "request2":
            jsonstring = "{a:1}";
            break;
    }
    var response = new SelectObjectContentResponse
    {
        Payload = new MockSelectObjectContentEventStream(
            new List<IS3Event> {
                new RecordsEvent(
                    new EventStreamMessage(
                        new List<IEventStreamHeader> { }, Encoding.UTF8.GetBytes(jsonstring)))}
            )
    };
    return response;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 09:02:32