如何用Moq在.NET中Mock AWS 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

