如何将大型对象序列化为JSON并流式写入SQL Server varbinary列
问题描述
我想要将一个大型对象序列化为JSON并存储到SQL Server的varbinary列中,理想情况下该序列化过程可直接分块写入数据库(流式处理)。
我编写了以下概念验证代码:
using System.Data; using System.Data.SqlClient; using System.Text; using System.Text.Json; Thread.Sleep(TimeSpan.FromSeconds(5)); var random = new Random(); var sb = new StringBuilder(); for (int i = 0; i < 125_000_000; i++) sb.Append(Convert.ToChar(random.NextInt64(26) + 97)); var sbAsString = sb.ToString(); Thread.Sleep(TimeSpan.FromSeconds(5)); SqlConnection sqlConnection = new("Server=localhost;Database=serializationtest;Trusted_Connection=True"); sqlConnection.Open(); using SqlCommand sqlCommand = new("insert into test (value) values (@value)", sqlConnection); var valueStream = new MemoryStream(); JsonSerializer.Serialize(valueStream, sbAsString); valueStream.Position = 0; SqlParameter parameter = sqlCommand.Parameters.Add("@value", SqlDbType.Binary, -1); parameter.Value = valueStream; sqlCommand.ExecuteNonQuery(); Thread.Sleep(TimeSpan.FromSeconds(5)); Console.WriteLine("Just for debugger");
显然这段代码会将序列化后的数据暂存于内存中,内存使用情况表现为两次明显的峰值:第一次是创建大型字符串时的内存增长,第二次是将JsonSerializer.Serialize结果写入MemoryStream时的内存增长。
请问是否可实现类似JsonSerializer.Serialize(BUFFERED_STREAM_TO_DB)的操作,直接写入数据库以减少第二次内存增长?
解决方案
完全可以实现流式序列化并直接写入数据库,避免将序列化后的JSON暂存到内存中。以下是两种可行方案:
方案一:使用SqlFileStream(需启用SQL Server FILESTREAM)
如果你的SQL Server启用了FILESTREAM功能,可以直接让JsonSerializer将序列化结果写入SqlFileStream,全程无需中间内存缓存:
- 前置条件:数据库已开启FILESTREAM,目标表的列类型设置为
varbinary(max) FILESTREAM。 - 代码示例:
using System.Data; using System.Data.SqlClient; using System.Text; using System.Text.Json; using System.IO; var random = new Random(); var sb = new StringBuilder(); for (int i = 0; i < 125_000_000; i++) sb.Append(Convert.ToChar(random.NextInt64(26) + 97)); var sbAsString = sb.ToString(); using var sqlConnection = new SqlConnection("Server=localhost;Database=serializationtest;Trusted_Connection=True"); sqlConnection.Open(); // 插入空记录并获取FILESTREAM访问路径 using var insertCmd = new SqlCommand(@" INSERT INTO test (value) VALUES (DEFAULT) OUTPUT GetFileStreamTransact-SQL(value, 1)", sqlConnection); var filePath = insertCmd.ExecuteScalar().ToString(); // 直接将JSON序列化结果写入SqlFileStream using var fileStream = new SqlFileStream(filePath, FileAccess.Write, FileOptions.SequentialScan, FileShare.None); JsonSerializer.Serialize(fileStream, sbAsString); fileStream.Close();
方案二:自定义流式写入Stream(无需FILESTREAM)
如果无法启用FILESTREAM,可以自定义Stream子类,将JSON序列化的输出分块提交到SQL Server,避免内存缓存:
- 自定义Stream实现:
public class SqlParameterStream : Stream { private readonly SqlCommand _command; private readonly SqlParameter _parameter; private readonly byte[] _buffer = new byte[8192]; // 8KB缓冲区,可根据需求调整 private int _bufferPosition; public SqlParameterStream(SqlCommand command, SqlParameter parameter) { _command = command; _parameter = parameter; _command.Prepare(); // 预编译命令提升性能 } public override bool CanRead => false; public override bool CanSeek => false; public override bool CanWrite => true; public override long Length => throw new NotSupportedException(); public override long Position { get => throw new NotSupportedException(); set => throw new NotSupportedException(); } public override void Flush() { if (_bufferPosition == 0) return; // 将当前缓冲区数据提交到数据库 using var ms = new MemoryStream(_buffer, 0, _bufferPosition); _parameter.Value = ms; _command.ExecuteNonQuery(); _bufferPosition = 0; } public override int Read(byte[] buffer, int offset, int count) => throw new NotSupportedException(); public override long Seek(long offset, SeekOrigin origin) => throw new NotSupportedException(); public override void SetLength(long value) => throw new NotSupportedException(); public override void Write(byte[] buffer, int offset, int count) { while (count > 0) { var copySize = Math.Min(count, _buffer.Length - _bufferPosition); Array.Copy(buffer, offset, _buffer, _bufferPosition, copySize); _bufferPosition += copySize; offset += copySize; count -= copySize; // 缓冲区满时提交数据 if (_bufferPosition == _buffer.Length) { Flush(); } } } protected override void Dispose(bool disposing) { if (disposing) { Flush(); // 处理剩余缓冲区数据 } base.Dispose(disposing); } }
- 使用自定义流的业务代码:
using System.Data; using System.Data.SqlClient; using System.Text; using System.Text.Json; var random = new Random(); var sb = new StringBuilder(); for (int i = 0; i < 125_000_000; i++) sb.Append(Convert.ToChar(random.NextInt64(26) + 97)); var sbAsString = sb.ToString(); using var sqlConnection = new SqlConnection("Server=localhost;Database=serializationtest;Trusted_Connection=True"); sqlConnection.Open(); // 建议包裹事务保证数据一致性 using var transaction = sqlConnection.BeginTransaction(); try { using var sqlCommand = new SqlCommand("INSERT INTO test (value) VALUES (@value)", sqlConnection, transaction); var parameter = sqlCommand.Parameters.Add("@value", SqlDbType.VarBinary, -1); using var sqlStream = new SqlParameterStream(sqlCommand, parameter); JsonSerializer.Serialize(sqlStream, sbAsString); transaction.Commit(); } catch { transaction.Rollback(); throw; }
额外优化建议
针对你代码中第一次内存增长(大型字符串创建),如果原始数据来自外部数据源(如文件、网络流),建议直接从数据源流式读取并序列化,避免一次性加载到内存中,进一步降低内存占用。
内容的提问来源于stack exchange,提问作者Wouter
相关产品推荐
相关产品推荐

