如何在C#中创建监听器从Amazon SQS获取10条消息并存入数据库?
解决C#中从Amazon SQS批量获取消息并存储到数据库的问题
看起来你遇到的核心问题不是没拿到10条消息,而是消息处理逻辑有误——你的代码把所有消息的Body拼接成一个字符串后只反序列化一次,这会破坏JSON格式(因为每条消息都是独立的ContactsDTO JSON对象),最终只能解析出第一条消息的内容,让你误以为只拿到了一条。
原代码的问题点
message += objMessage.Body;将多条独立的JSON对象拼接成了无效的字符串,比如变成{"FirstName":"A"...}{"FirstName":"B"...},这种格式无法被正确反序列化为多个ContactsDTO对象。- 仅执行一次反序列化,只能得到第一个有效JSON对象的内容,后续消息的内容会被忽略或导致反序列化错误。
修正后的完整实现
下面是调整后的代码,会批量获取消息、逐个解析并存储到数据库,同时处理消息删除避免重复消费:
using System.Collections.Generic; using System.Web.Script.Serialization; // 注意:如果是.NET Core/.NET 5+,建议使用System.Text.Json替代JavaScriptSerializer using Amazon.SQS; using Amazon.SQS.Model; // 初始化SQS客户端(确保你已经正确配置了AWS凭证) var objClient = new AmazonSQSClient(); // 构建获取消息的请求 var receiveMessageRequest = new ReceiveMessageRequest { QueueUrl = "你的队列URL", MaxNumberOfMessages = 10, // 最多获取10条消息 WaitTimeSeconds = 20, // 开启长轮询,减少空轮询次数,最多20秒 VisibilityTimeout = 30 // 设置消息可见性超时,避免处理期间被其他消费者获取 }; // 获取消息响应 var receiveMessageResponse = objClient.ReceiveMessage(receiveMessageRequest); var contactsToSave = new List<ContactsDTO>(); var receiptHandles = new List<string>(); // 保存已处理消息的句柄,用于后续删除 // 初始化序列化器(.NET Core/.NET 5+推荐使用JsonSerializer) var serializer = new JavaScriptSerializer(); foreach (var message in receiveMessageResponse.Messages) { try { // 逐个解析每条消息的Body var contact = serializer.Deserialize<ContactsDTO>(message.Body); contactsToSave.Add(contact); receiptHandles.Add(message.ReceiptHandle); // 也可以在这里直接单条存入数据库,不需要先收集到列表 // SaveSingleContactToDatabase(contact); } catch (System.Exception ex) { // 处理解析失败的情况,避免一条消息出错导致批量处理中断 System.Console.WriteLine($"解析消息失败: {ex.Message},消息内容: {message.Body}"); } } // 批量将解析后的联系人存入数据库 if (contactsToSave.Count > 0) { SaveContactsToDatabase(contactsToSave); System.Console.WriteLine($"成功处理并存储{contactsToSave.Count}条消息"); } // 批量删除已处理的消息,避免重复消费 if (receiptHandles.Count > 0) { var deleteBatchRequest = new DeleteMessageBatchRequest { QueueUrl = "你的队列URL", Entries = receiptHandles.Select((handle, index) => new DeleteMessageBatchRequestEntry { Id = index.ToString(), // 每个条目需要唯一ID ReceiptHandle = handle }).ToList() }; var deleteResponse = objClient.DeleteMessageBatch(deleteBatchRequest); // 检查是否有删除失败的消息 foreach (var failedEntry in deleteResponse.Failed) { System.Console.WriteLine($"删除消息失败ID {failedEntry.Id}: {failedEntry.Message}"); } } // 数据库批量存储的示例方法(根据你的数据库框架调整,比如EF Core/ADO.NET) void SaveContactsToDatabase(List<ContactsDTO> contacts) { // 示例:使用EF Core using (var dbContext = new YourDatabaseContext()) { dbContext.Contacts.AddRange(contacts); dbContext.SaveChanges(); } // 如果你用ADO.NET,可以编写批量插入的SQL语句 // using (var conn = new SqlConnection("你的数据库连接字符串")) // { // conn.Open(); // // 执行批量插入逻辑 // } } // 你的ContactsDTO类保持不变 public class ContactsDTO { public string FirstName { get; set; } public string LastName { get; set; } public string Address { get; set; } }
额外注意事项
- 关于
MaxNumberOfMessages:SQS最多返回10条消息,但实际返回数量取决于队列中的消息总数、消息大小等因素,如果队列里只有3条,就只会返回3条。 - 序列化器选择:
JavaScriptSerializer是旧的.NET Framework组件,如果你使用的是.NET Core/.NET 5+,建议使用System.Text.Json.JsonSerializer或者Newtonsoft.Json(Json.NET),它们的性能和兼容性更好。 - 消息可见性超时:设置
VisibilityTimeout可以确保在你处理消息的这段时间内,其他消费者不会获取到同一条消息,避免重复处理。 - 错误处理:一定要添加异常捕获,避免单个消息的格式错误导致整个批量处理流程中断。
内容的提问来源于stack exchange,提问作者Rajat
相关产品推荐
相关产品推荐

