解决MySql布尔型新旧版本兼容问题,采用枚举来表示布尔型的数据表。由正向工程赋值
大石头 authored at 2018-05-15 21:21:05
13.05 KiB
X
using BenchmarkDotNet.Attributes;
using NewLife;
using NewLife.Net;
using System.Net.Sockets;

namespace Benchmark.NetBenchmarks;

/// <summary>网络库服务端接收吞吐量基准测试</summary>
/// <remarks>
/// 包含两个测试方法,同表对比:
/// 1. 逐包发送:每次 Send(32B),衡量每次 recv() 回调的完整处理开销
/// 2. 批量发送:256 包合并为 Send(8KB),衡量 TCP 流式吞吐(开销被 TCP 粘包分摊)
/// 命令:dotnet run --project Benchmark/Benchmark.csproj -c Release -- --filter "*NetServerThroughputBenchmark*"
/// </remarks>
[MemoryDiagnoser]
[GcServer(true)]
[SimpleJob(warmupCount: 2, iterationCount: 5)]
public class NetServerThroughputBenchmark : IDisposable
{
    /// <summary>逐包发送逻辑包总数(2^21),C=1 时迭代约 1 秒</summary>
    private const Int32 PerPacketTotal = 2_097_152;

    /// <summary>批量发送逻辑包总数(2^24),确保迭代 >100ms</summary>
    private const Int32 BatchTotal = 16_777_216;

    /// <summary>批量发送合并数:256 个 32B 包 = 8KB 一次 Send</summary>
    private const Int32 BatchSize = 256;

    private const Int32 Port = 7779;

    private ThroughputNetServer? _server;
    private ISocketClient[] _clients = null!;
    private Byte[] _singlePayload = null!;
    private Byte[] _batchPayload = null!;

    /// <summary>数据包大小(字节)</summary>
    [Params(32)]
    public Int32 PacketSize { get; set; }

    /// <summary>并发客户端数</summary>
    [Params(1, 4, 16, 64, 256, 1024)]
    public Int32 Concurrency { get; set; }

    /// <summary>全局初始化:启动服务端,建立所有客户端连接</summary>
    [GlobalSetup]
    public void Setup()
    {
        _singlePayload = new Byte[PacketSize];
        Random.Shared.NextBytes(_singlePayload);

        _batchPayload = new Byte[BatchSize * PacketSize];
        for (var i = 0; i < BatchSize; i++)
            Buffer.BlockCopy(_singlePayload, 0, _batchPayload, i * PacketSize, PacketSize);

        // 增大 IOCP 接收缓冲区
        SocketSetting.Current.BufferSize = 64 * 1024;

        _server = new ThroughputNetServer
        {
            Port = Port,
            ProtocolType = NetType.Tcp,
            AddressFamily = AddressFamily.InterNetwork,
            UseSession = false,
        };
        _server.Start();

        _clients = new ISocketClient[Concurrency];
        for (var i = 0; i < Concurrency; i++)
        {
            var client = new NetUri($"tcp://127.0.0.1:{Port}").CreateRemote();
            client.Open();
            _clients[i] = client;
        }
    }

    /// <summary>逐包发送:每次 Send(32B),测量每次 recv() 回调的完整处理开销</summary>
    [Benchmark(Description = "逐包发送", OperationsPerInvoke = PerPacketTotal)]
    public Int64 PerPacketThroughput()
    {
        _server!.Reset((Int64)PacketSize * PerPacketTotal);

        var perClient = PerPacketTotal / Concurrency;
        var tasks = new Task[Concurrency];
        for (var i = 0; i < Concurrency; i++)
        {
            var idx = i;
            tasks[i] = Task.Run(() =>
            {
                var client = _clients[idx];
                var payload = _singlePayload;
                for (var n = 0; n < perClient; n++)
                    client.Send(payload);
            });
        }

        Task.WaitAll(tasks);
        if (!_server!.WaitComplete(120_000))
            throw new TimeoutException($"逐包发送超时(已收 {_server.ReceivedBytes:N0} / {(Int64)PacketSize * PerPacketTotal:N0})");

        return _server.ReceivedBytes;
    }

    /// <summary>批量发送:256 包合并为 8KB 一次 Send,测量 TCP 流式吞吐上限</summary>
    [Benchmark(Description = "批量发送", OperationsPerInvoke = BatchTotal)]
    public Int64 BatchThroughput()
    {
        _server!.Reset((Int64)PacketSize * BatchTotal);

        var perClient = BatchTotal / Concurrency;
        var sendsPerClient = perClient / BatchSize;
        var tasks = new Task[Concurrency];
        for (var i = 0; i < Concurrency; i++)
        {
            var idx = i;
            tasks[i] = Task.Run(() =>
            {
                var client = _clients[idx];
                var payload = _batchPayload;
                for (var n = 0; n < sendsPerClient; n++)
                    client.Send(payload);
            });
        }

        Task.WaitAll(tasks);
        if (!_server!.WaitComplete(60_000))
            throw new TimeoutException($"批量发送超时(已收 {_server.ReceivedBytes:N0} / {(Int64)PacketSize * BatchTotal:N0})");

        return _server.ReceivedBytes;
    }

    /// <summary>全局清理:释放所有客户端和服务端</summary>
    [GlobalCleanup]
    public void Cleanup()
    {
        if (_clients != null)
        {
            foreach (var c in _clients)
                c?.Dispose();
            _clients = null!;
        }

        _server?.Dispose();
        _server = null;
    }

    /// <summary>释放资源</summary>
    public void Dispose() => Cleanup();
}

/// <summary>服务端仅接收计数,不回发,专用于吞吐量测试</summary>
class ThroughputNetServer : NetServer
{
    private Int64 _receivedBytes;
    private Int64 _expectedBytes;
    private readonly ManualResetEventSlim _completed = new(false);

    /// <summary>已接收总字节数</summary>
    public Int64 ReceivedBytes => Interlocked.Read(ref _receivedBytes);

    /// <summary>重置计数器并设置期望接收字节数</summary>
    /// <param name="expectedBytes">本轮期望接收的总字节数</param>
    public void Reset(Int64 expectedBytes)
    {
        Interlocked.Exchange(ref _receivedBytes, 0);
        Interlocked.Exchange(ref _expectedBytes, expectedBytes);
        _completed.Reset();
    }

    /// <summary>等待服务端接收完成</summary>
    /// <param name="millisecondsTimeout">超时毫秒数</param>
    /// <returns>是否在超时前接收完成</returns>
    public Boolean WaitComplete(Int32 millisecondsTimeout) => _completed.Wait(millisecondsTimeout);

    /// <summary>接收数据仅计数,不回发</summary>
    protected override void OnReceive(INetSession session, ReceivedEventArgs e)
    {
        var bytes = e.Packet?.Total ?? 0;
        if (bytes <= 0) return;

        var expected = Interlocked.Read(ref _expectedBytes);
        var total = Interlocked.Add(ref _receivedBytes, bytes);
        if (total >= expected)
            _completed.Set();
    }
}

/// <summary>裸 Socket 层 TCP 回声基准:逐包往返与流水线两种模式,覆盖包大小与并发维度</summary>
/// <remarks>
/// <para>对比 NetServerThroughputBenchmark(单向接收、32B 固定维度):本类走完整回声链路(收→回发→收),
/// 维度为 包大小 × 并发,可用于评估不同报文规模下的往返链路成本与流水线吞吐上限。</para>
/// <para>换算说明:单迭代数据量按 <see cref="TargetBytes"/> 字节限流,小包档以往返次数/包数上限保护;
/// 不设 OperationsPerInvoke,报表 ns/op 即“每迭代”成本;
/// MB/s = 传输字节 ÷ 迭代耗时,msg/s = 包数 ÷ 迭代耗时。</para>
/// <para>命令:dotnet run --project Benchmark/Benchmark.csproj -c Release -- --filter "*NetEchoBenchmark*"</para>
/// </remarks>
[MemoryDiagnoser]
[GcServer(true)]
[SimpleJob(warmupCount: 2, iterationCount: 5)]
public class NetEchoBenchmark : IDisposable
{
    /// <summary>单迭代数据传输字节目标(大包档按字节限流)</summary>
    private const Int64 TargetBytes = 256L * 1024 * 1024;

    /// <summary>往返模式单迭代总往返次数上限(小包档防迭代过长)</summary>
    private const Int32 MaxRoundTrips = 50_000;

    /// <summary>流水线模式单迭代总包数上限(小包档防迭代过长)</summary>
    private const Int32 MaxPackets = 2_000_000;

    private const Int32 EchoPort = 7780;

    /// <summary>接收兜底超时,防基准迭代永久挂起</summary>
    private const Int32 EchoTimeoutMs = 120_000;

    private NetServer? _echoServer;
    private ISocketClient[] _echoClients = null!;
    private Byte[] _echoPayload = null!;

    /// <summary>数据包大小(字节)</summary>
    [Params(16, 256, 4096, 65536, 1048576)]
    public Int32 EchoPacketSize { get; set; }

    /// <summary>并发客户端数</summary>
    [Params(1, 4, 16, 64)]
    public Int32 EchoConcurrency { get; set; }

    /// <summary>每客户端单迭代往返数(字节目标与次数上限取小者)</summary>
    private Int32 RoundsPerClient => Math.Max(1, (Int32)Math.Min(TargetBytes / EchoPacketSize, MaxRoundTrips) / EchoConcurrency);

    /// <summary>每客户端单迭代包数(字节目标与包数上限取小者)</summary>
    private Int32 PacketsPerClient => Math.Max(1, (Int32)Math.Min(TargetBytes / EchoPacketSize, MaxPackets) / EchoConcurrency);

    /// <summary>全局初始化:启动回声服务端并建立全部客户端</summary>
    [GlobalSetup(Targets = [nameof(RoundTrip), nameof(Pipeline)])]
    public void EchoSetup()
    {
        // 高并发同步往返会触发线程池注入限速(~1线程/秒),预热最小线程避免秒级毛刺污染测量
        ThreadPool.SetMinThreads(512, 512);

        _echoPayload = new Byte[EchoPacketSize];
        Random.Shared.NextBytes(_echoPayload);

        // 增大接收缓冲区,降低大包分段的系统调用开销
        SocketSetting.Current.BufferSize = 256 * 1024;

        _echoServer = new NetServer
        {
            Port = EchoPort,
            ProtocolType = NetType.Tcp,
            AddressFamily = AddressFamily.InterNetwork,
        };
        _echoServer.Received += (s, e) =>
        {
            if (e.Packet != null && s is INetSession session) session.Send(e.Packet);
        };
        _echoServer.Start();

        _echoClients = new ISocketClient[EchoConcurrency];
        for (var i = 0; i < EchoConcurrency; i++)
        {
            // 拉取模式:每客户端独享直读回显,避免事件分发干扰基准
            var client = new TcpSession
            {
                Remote = new NetUri($"tcp://127.0.0.1:{EchoPort}"),
                AutoReceive = false,
                BufferSize = Math.Max(64 * 1024, EchoPacketSize),
                Timeout = EchoTimeoutMs,
            };
            client.Open();
            _echoClients[i] = client;
        }
    }

    /// <summary>逐包往返:每客户端发一包读满一包回显,测最小往返链路开销</summary>
    [Benchmark(Description = "逐包往返")]
    public Int64 RoundTrip()
    {
        var rounds = RoundsPerClient;
        var tasks = new Task[EchoConcurrency];
        for (var i = 0; i < EchoConcurrency; i++)
        {
            var idx = i;
            tasks[i] = Task.Run(() => EchoLoop(_echoClients[idx], rounds));
        }

        Task.WaitAll(tasks);

        return (Int64)EchoPacketSize * rounds * EchoConcurrency;
    }

    /// <summary>流水线:各客户端持续发送不等待回显,服务端边收边回,客户端按字节读满</summary>
    [Benchmark(Description = "流水线")]
    public Int64 Pipeline()
    {
        var count = PacketsPerClient;
        var tasks = new Task[EchoConcurrency];
        for (var i = 0; i < EchoConcurrency; i++)
        {
            var idx = i;
            tasks[i] = Task.Run(() => SendAndReceive(_echoClients[idx], count));
        }

        Task.WaitAll(tasks);

        return (Int64)EchoPacketSize * count * EchoConcurrency;
    }

    /// <summary>单客户端逐包往返:发一包后读满同字节回显,再发下一包</summary>
    private void EchoLoop(ISocketClient client, Int32 rounds)
    {
        for (var n = 0; n < rounds; n++)
        {
            client.Send(_echoPayload);

            var received = 0;
            while (received < EchoPacketSize)
            {
                using var pk = client.Receive();
                if (pk == null || pk.Length <= 0) throw new TimeoutException("回显接收中断");

                received += pk.Length;
            }
        }
    }

    /// <summary>单客户端流水线:发送与回显读取并行,按字节计数读满目标流量</summary>
    private void SendAndReceive(ISocketClient client, Int32 count)
    {
        var expected = (Int64)EchoPacketSize * count;
        var received = 0L;

        var sendTask = Task.Run(() =>
        {
            for (var n = 0; n < count; n++)
                client.Send(_echoPayload);
        });

        while (received < expected)
        {
            using var pk = client.Receive();
            if (pk == null || pk.Length <= 0) throw new TimeoutException("回显接收中断");

            received += pk.Length;
        }

        sendTask.Wait();
    }

    /// <summary>回声基准清理:释放所有客户端和服务端</summary>
    [GlobalCleanup(Targets = [nameof(RoundTrip), nameof(Pipeline)])]
    public void EchoCleanup()
    {
        if (_echoClients != null)
        {
            foreach (var c in _echoClients)
                c?.Dispose();
            _echoClients = null!;
        }

        _echoServer?.Dispose();
        _echoServer = null;
    }

    /// <summary>释放资源</summary>
    public void Dispose() => EchoCleanup();
}