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

C#应用从SQL Server取海量数据转Stream推送至Storage的技术问询

Solution for Streaming Large SQL Results to Storage App

Why Entity Framework Isn't Ideal Here

EF's Database.SqlQuery<T> requires a predefined type T to map query results, which isn't feasible when your query's output schema is dynamic (since @insertedParams can be anything). Additionally, even if you use dynamic, EF may buffer results in memory—something that's impossible for TB-sized datasets.

Use ADO.NET's SqlDataReader with streaming capabilities to process rows incrementally, avoiding loading the entire dataset into memory. Combine this with a pipe to stream data directly to the Storage app's PushData method.

Step-by-Step Implementation

Here's an async implementation that handles dynamic query results and streams data efficiently:

using System.Data;
using System.Data.SqlClient;
using System.IO;
using System.IO.Pipelines;
using System.Linq;
using System.Threading.Tasks;

public async Task<Guid> StreamQueryResultsToStorage(string connectionString, string sqlRequest)
{
    using var connection = new SqlConnection(connectionString);
    await connection.OpenAsync();

    using var command = new SqlCommand(sqlRequest, connection);
    // SequentialAccess enables efficient streaming of large columns
    using var reader = await command.ExecuteReaderAsync(CommandBehavior.SequentialAccess);

    // Retrieve column names from the query result schema
    var schemaTable = reader.GetSchemaTable();
    var columnNames = schemaTable?
        .Rows.Cast<DataRow>()
        .Select(row => row["ColumnName"].ToString())
        .ToList() ?? new List<string>();

    // Use Pipe to stream data between reader and Storage app
    var pipe = new Pipe();
    
    // Background task to write query results to the pipe
    var writeTask = Task.Run(async () =>
    {
        await using var streamWriter = new StreamWriter(pipe.Writer.AsStream());
        
        // Write CSV header (adjust format based on Storage app's requirements)
        await streamWriter.WriteLineAsync(string.Join(",", columnNames));
        
        while (await reader.ReadAsync())
        {
            var rowValues = new List<string>();
            for (int colIndex = 0; colIndex < reader.FieldCount; colIndex++)
            {
                if (reader.IsDBNull(colIndex))
                {
                    rowValues.Add("NULL");
                    continue;
                }

                var fieldType = reader.GetFieldType(colIndex);
                string value;

                // Handle large binary columns (e.g., VARBINARY(MAX))
                if (fieldType == typeof(byte[]))
                {
                    var buffer = new byte[4096];
                    long bytesRead;
                    long offset = 0;
                    using var ms = new MemoryStream();
                    while ((bytesRead = await reader.GetBytesAsync(colIndex, offset, buffer, 0, buffer.Length)) > 0)
                    {
                        ms.Write(buffer, 0, (int)bytesRead);
                        offset += bytesRead;
                    }
                    value = Convert.ToBase64String(ms.ToArray());
                }
                // Handle other data types
                else
                {
                    value = reader.GetValue(colIndex).ToString() ?? string.Empty;
                    // Escape CSV special characters
                    value = $"\"{value.Replace("\"", "\"\"")}\"";
                }

                rowValues.Add(value);
            }

            await streamWriter.WriteLineAsync(string.Join(",", rowValues));
            await streamWriter.FlushAsync();
        }

        // Signal the pipe that writing is complete
        await pipe.Writer.CompleteAsync();
    });

    // Pass the pipe's read stream to the Storage app
    var storageGuid = Storage.PushData(pipe.Reader.AsStream());

    // Wait for the write task to finish to ensure all data is sent
    await writeTask;

    return storageGuid;
}

Key Details

  • CommandBehavior.SequentialAccess: Critical for handling large columns (like VARCHAR(MAX) or VARBINARY(MAX)) by allowing incremental reading instead of loading entire columns into memory.
  • Dynamic Schema Handling: GetSchemaTable() retrieves column metadata, so you don't need to know the query's output structure in advance.
  • Pipe Streaming: The Pipe class from System.IO.Pipelines enables efficient, memory-safe streaming between the SQL reader and the Storage app. Data is written incrementally as it's read, avoiding in-memory buffering of the entire dataset.
  • Async Execution: Uses async methods to avoid blocking threads, improving scalability for concurrent requests.

Adjustments for Your Use Case

  • Data Format: Replace the CSV writing logic with the format expected by your Storage app (e.g., JSON lines, binary, or a custom format).
  • Error Handling: Add try/catch blocks to handle exceptions from database operations, streaming, or the Storage app.
  • Large Text Columns: For VARCHAR(MAX) columns, use reader.GetChars() instead of GetValue() to read in chunks, similar to the binary column handling.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 21:09:30