引言

在实际项目中,我们经常需要处理来自不同厂商、不同格式的日志文件——华为用 MML 命令格式写配置,中兴用专有 XML,爱立信又有自己的一套。如果每种格式都写一个独立的解析脚本,代码会变成一堆 if/else 的巨型分支,维护起来痛苦不堪。

SwitchData 项目面临的就是这样的场景:需要解析华为 PGW 的 SAAC 格式日志、华为 DNS 的 Zone 配置、通用 MML 命令日志等至少四种格式。最终我们落地了一套工厂策略模式 + Channel 异步管道 + 批量持久化的完整架构,解耦、高效、易扩展。

本文就带你从零拆解这套架构,看它是如何用约 500 行核心代码支撑起整个网元日志解析流水线的。

一、核心设计:接口抽象 + 工厂分发

架构的第一步永远是定义不变量。无论厂商和格式怎么变,“解析日志 → 入库”这个流程是不变的。我们把这个不变的契约抽成接口 INeLogParser:

public interface INeLogParser
{
    int BatchSize { get; }
    string FormatCode { get; }
    
    Task ParseAndPersistAsync(
        NeDeviceParseContext context,
        NeDeviceParseResult parseResult,
        CancellationToken cancellationToken);
}

三个成员各司其职:

成员 作用
BatchSize 批量入库的行数阈值,不同解析器可按需调整
FormatCode 格式标识符,工厂用它来路由到具体实现
ParseAndPersistAsync 解析并持久化的主入口,异步流式处理

然后用静态工厂按 FormatCode 分发:

public static class NeLogParserFactory
{
    private static readonly Dictionary<string, INeLogParser> _parsers;

    static NeLogParserFactory()
    {
        _parsers = new(StringComparer.OrdinalIgnoreCase)
        {
            ["MML"]  = new MmlLogParser(),
            ["SAAC"] = new HwPgwParser(),
            ["HW_DNS"] = new HwDnsParser(),
            ["HW_DNS_LOG"] = new HwDnsLogParser()
        };
    }

    public static bool TryGetParser(string formatCode, out INeLogParser parser)
        => _parsers.TryGetValue(formatCode, out parser);
}

工厂 + 策略的组合在这里发挥了威力:上层业务代码只需要知道 formatCode,不需要知道背后有多少种解析器、各自内部实现如何。新增一种日志格式时,写一个新的 INeLogParser 实现,在工厂里注册一行就完事,调用方零改动。

二、Channel 异步管道:解析与持久化解耦

传统的解析代码通常是这样的:读一行 → 解析 → 立刻写数据库 → 再读下一行。这种”单线程串行”模式下,磁盘 IO 和数据库 IO 互相阻塞,整体吞吐量上不去。

Channel 解决了这个问题。它是 .NET 6 引入的异步生产者-消费者队列,天生适配”读-解析”和”批量入库”两个独立阶段:

graph LR
    A[日志文件] --> B[Channel.Writer 生产者<br/>逐行解析]
    B --> C["Channel 有界队列<br/>BatchSize x 2"]
    C --> D[Channel.Reader 消费者<br/>await foreach 批量入库]

以 MmlLogParser 为例,核心骨架只有几十行:

public async Task ParseAndPersistAsync(
    NeDeviceParseContext parseContext,
    NeDeviceParseResult parseResult,
    CancellationToken cancellationToken)
{
    // 1. 创建有界 Channel,容量 = BatchSize * 2
    var channel = Channel.CreateBounded<MmlCommand>(
        new(BatchSize * 2)
        {
            FullMode = BoundedChannelFullMode.Wait,
            SingleReader = true,
            SingleWriter = true
        });

    // 2. 消费者:启动后立刻阻塞等待,不浪费 CPU
    Task consumerTask = PersistAsync(channel.Reader, parseContext, parseResult, cancellationToken);

    try
    {
        // 3. 生产者:逐行读取 → 解析 → 写入 Channel
        await ParseAsync(parseContext, channel.Writer, cancellationToken);
    }
    finally
    {
        // 4. 标记写完,让消费者能正常退出
        channel.Writer.Complete();
        await consumerTask; // 等消费者把积压的数据刷完
    }
}

几个关键设计点:

有界 + Wait 背压:FullMode.Wait 意味着当队列满时,生产者会异步等待,而不是丢弃数据或阻塞线程。这天然实现了消费者限速——数据库忙的时候,文件读取自然变慢,不会把内存撑爆。

单读单写:SingleReader = true, SingleWriter = true 让 Channel 内部可以跳过锁竞争,在这个解析场景下(一个文件 → 一个 Channel → 一个消费者)是安全的。

try/finally 保证完成:channel.Writer.Complete() 放在 finally 里,确保即使解析中途抛异常,消费者也能收到”结束”信号,正常刷完已解析的数据再退出,不会丢数据。

三、两种解析器实现,看策略差异

通用 MML 解析器

MML(Man Machine Language)是华为/中兴通用的命令格式,形如 LST PCCSUBDATA:APN="iot";。MmlLogParser 用正则 ^[A-Z]+\s+.+?;$ 匹配命令行,逐条解析后灌入 Channel:

private static async Task ParseAsync(
    NeDeviceParseContext parseContext,
    ChannelWriter<MmlCommand> writer,
    CancellationToken cancellationToken)
{
    using var reader = new StreamReader(parseContext.NeLogFile.FullPath);

    while (await reader.ReadLineAsync(cancellationToken) is { } line)
    {
        line = line.Trim();
        if (line.Length == 0 || !MmlCommandRegex.IsMatch(line)) continue;

        var cmd = new MmlCommand(line);
        // 过滤掉不关心的对象类型
        if (!objectSet.Contains(cmd.ObjectName)) continue;

        await writer.WriteAsync(cmd, cancellationToken);
    }
}

特点:纯文本 + 正则,BatchSize = 10000,通用但解析精度依赖正则匹配。

华为 PGW SAAC 解析器

HwPgwParser 处理的是华为 PGW 的专有配置格式,结构更规范,每种业务对象(Host、Filter、Rule、UserProfile…)有独立的 C# 类:

private static readonly Dictionary<string, Type> BusinessTableTypes = new(StringComparer.OrdinalIgnoreCase)
{
    { HostEntry.BusinessTableName, typeof(HostEntry) },
    { FilterEntry.BusinessTableName, typeof(FilterEntry) },
    { RuleEntry.BusinessTableName, typeof(RuleEntry) },
    { UserProfile.BusinessTableName, typeof(UserProfile) },
    // ... 共 13 种
};

特点:强类型 + 映射表,BatchSize = 20000,解析速度快(不需要正则),但需要为每种格式维护对应的 Entry 类。

两个解析器共享同一个 Channel 管道模式,只是生产者的解析逻辑不同——这就是策略模式的威力。

四、批量 Flush + ADO.NET BulkCopy:持久化加速

消费者侧由 NeLogPersistService.PersistAsync 统一处理,核心是 BatchPersistContext:

private sealed class BatchPersistContext
{
    public required Type RecordType { get; init; }
    public required DataTable DataTable { get; init; }
    public required NeTableParseResult Result { get; init; }
}

用 Dictionary<Type, BatchPersistContext> 按记录类型分桶,每桶一个独立的 DataTable,达到 BatchSize 阈值时用 ADO.NET BulkCopy 批量写入:

await foreach (INeLogRecord record in reader.ReadAllAsync(cancellationToken))
{
    var persistContext = batchPersistContexts[record.GetType()];
    DataTableConverter.AddRow(persistContext.DataTable, record);

    await FlushAsync(persistContext, force: false, batchSize, cancellationToken);
}
// 循环结束后,force=true 把残余数据也刷掉
foreach (var ctx in batchPersistContexts.Values)
    await FlushAsync(ctx, force: true, batchSize, cancellationToken);

FlushAsync 的判定逻辑很简洁:

if (!force && dataTable.Rows.Count < batchSize)
    return; // 不够一批,继续攒

await DbHelper.SwitchNeData.Database.BulkCopyAsync(
    dataTable, dbTableName, batchSize, addColumnMapping: true);
dataTable.Clear(); // 清空攒下一批

为什么 BulkCopy 这么快?因为它绕过了 INSERT 语句的解析、参数化、执行计划等开销,直接把 DataTable 的内存块推给数据库。实测在百万级记录场景下,比逐条 INSERT 快 10~50 倍。

五、可扩展性验证:新增一种格式有多简单

假设现在要接入爱立信的 XML 配置日志 ERICSSON_XML,只需要三步:

第一步:实现解析器

public sealed class EricssonXmlParser : INeLogParser
{
    public int BatchSize => 15000;
    public string FormatCode => "ERICSSON_XML";

    public async Task ParseAndPersistAsync(...)
    {
        var channel = Channel.CreateBounded<EricssonRecord>(BatchSize * 2);
        // 生产者:用 XmlReader 流式解析 → 写入 Channel
        // 消费者:复用 NeLogPersistService 的批量写入逻辑
    }
}

第二步:注册到工厂

// NeLogParserFactory 的静态构造函数里加一行
_parsers["ERICSSON_XML"] = new EricssonXmlParser();

就这两处改动。上层业务调用方、调度系统、数据库表结构……全部不受影响。这就是开闭原则(OCP)的完美体现:对扩展开放,对修改关闭。

六、性能与稳定性设计总结

回顾整个架构,能稳定跑在生产环境,靠的是这几个关键选择:

设计点 解决什么问题 技术选型
工厂 + 策略 多格式分发、解耦调用方 NeLogParserFactory + INeLogParser
Channel 管道 IO 并行、背压控制 Channel.CreateBounded + await foreach
批量 Flush 减少数据库交互次数 BatchPersistContext + 阈值判定
ADO.NET BulkCopy 百万级快速入库 内存直推数据库
CancellationToken 长时间任务可取消 全链路透传
Try/Finally 异常时保证数据不丢 Writer.Complete() 保底

结语

这套架构已经在 SwitchData 项目生产环境稳定运行两年多,从最初的两种日志格式扩展到现在的四种,调用方代码一行没改过。它的核心思路其实很朴素:把变化的东西(解析逻辑)藏在策略后面,把不变的东西(管道骨架)固定下来,然后用 Channel 把两个阶段拼起来。

如果你正在设计类似的多格式数据处理系统,可以直接套用这个套路——哪怕不是日志解析,是电商订单、风控事件、IoT 遥测数据,同样适用。