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

如何将大型对象序列化为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,全程无需中间内存缓存:

  1. 前置条件:数据库已开启FILESTREAM,目标表的列类型设置为varbinary(max) FILESTREAM。
  2. 代码示例:
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,避免内存缓存:

  1. 自定义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);
    }
}
  1. 使用自定义流的业务代码:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 12:35:58