极简架构在数据平台中的复盘:ETL 管道的设计模式与容错机制
极简架构在数据平台中的复盘ETL 管道的设计模式与容错机制一、引言数据平台是过度工程化的重灾区。一个简单的 ETL 任务方案评审时往往被未来可能需要的需求堆成一个重型的分布式调度系统。去年做的一个内部数据报表平台我从零搭建了它的 ETL 管道刻意保持了极简架构。上线一年后回头看这个选择是正确的——系统稳定运行、代码量不到 3000 行、维护成本极低。这个平台的功能很简单每天凌晨从三个数据源MySQL 业务库、MongoDB 日志库、第三方 API拉取数据经过清洗和聚合后写入 ClickHouse供 BI 看板查询。日均处理数据量约 200 万行不算大规模但对可靠性和可维护性有硬性要求。本文将复盘 ETL 管道的架构设计、设计模式应用和容错机制的演进。二、架构设计三个核心原则原则一一个任务只做一件事ETL 管道的每个节点只负责一个明确的数据转换操作。不做一锅炖式的多源 Join也不在一个任务里混合 Extract 和 Load。// 管道节点的统一接口 type Stage interface { Name() string Process(ctx context.Context, input -chan *Record, output chan- *Record) error } // 管道编排器 type Pipeline struct { stages []Stage logger *slog.Logger } func (p *Pipeline) Run(ctx context.Context, reader RecordReader) error { var input chan *Record make(chan *Record, 1000) var output chan *Record // 启动数据读取 errCh : make(chan error, 1) go func() { defer close(input) errCh - reader.Read(ctx, input) }() // 串联每个 Stage for i, stage : range p.stages { output make(chan *Record, 1000) go func(s Stage, in, out chan *Record) { defer close(out) if err : s.Process(ctx, in, out); err ! nil { p.logger.Error(stage failed, slog.String(stage, s.Name()), slog.String(error, err.Error()), ) } }(stage, input, output) input output } // 最后的 output 就是最终结果写入目标 // ...写入逻辑... return -errCh }原则二配置驱动代码不变不同数据源的表结构、清洗规则、聚合逻辑全部放在 YAML 配置里。新增一个数据源只需要加一个配置文件不需要改一行 Go 代码。# configs/etl/order_daily.yaml source: type: mysql connection: ${MYSQL_DSN} query: | SELECT order_id, user_id, amount, status, created_at FROM orders WHERE created_at {start_date} AND created_at {end_date} stages: - name: clean_amount type: transform rules: - field: amount action: default_zero - field: status action: uppercase - name: enrich_user_tier type: lookup source: redis key: user_id field: user_tier default: unknown - name: aggregate_daily type: aggregate group_by: [user_tier, status] metrics: total_amount: SUM(amount) order_count: COUNT(order_id) avg_amount: AVG(amount) target: type: clickhouse table: dwd_order_daily mode: append配置驱动让非开发人员数据分析师也能新增和修改 ETL 任务极大降低了沟通成本。原则三有状态的管道是无状态的组合每个 Stage 本身可以是有状态的比如 Join 操作需要内存缓存但管道编排层只关心 Stage 的输入输出 Channel。这种设计让每个 Stage 可以独立测试、独立替换。三、设计模式三件套模式一Filter过滤器链脏数据过滤不适合做成一个大函数更适合用过滤器链——每个 Filter 独立判断一条记录是否需要丢弃。type Filter interface { Accept(record *Record) bool Reason() string } type FilterChain struct { filters []Filter } func (fc *FilterChain) Apply(records []*Record) ([]*Record, []FilteredRecord) { var passed []*Record var filtered []FilteredRecord for _, rec : range records { rejected : false for _, f : range fc.filters { if !f.Accept(rec) { filtered append(filtered, FilteredRecord{ Record: rec, Reason: f.Reason(), }) rejected true break } } if !rejected { passed append(passed, rec) } } return passed, filtered }模式二Strategy策略模式同一个字段的清洗规则因数据源不同而异。使用策略模式每个字段绑定一个清洗策略。type CleanStrategy interface { Clean(value interface{}) (interface{}, error) } type DefaultZeroStrategy struct{} func (s DefaultZeroStrategy) Clean(value interface{}) (interface{}, error) { if value nil { return 0, nil } return value, nil } type TrimStringStrategy struct{} func (s TrimStringStrategy) Clean(value interface{}) (interface{}, error) { if str, ok : value.(string); ok { return strings.TrimSpace(str), nil } return value, nil } var strategyRegistry map[string]CleanStrategy{ default_zero: DefaultZeroStrategy{}, trim_string: TrimStringStrategy{}, uppercase: UppercaseStrategy{}, }模式三Observer观察者模式管道运行过程中需要向外部暴露运行状态、进度、错误。用观察者模式解耦管道逻辑和监控逻辑。type PipelineObserver interface { OnStageStart(stageName string) OnStageProgress(stageName string, processed, total int64) OnStageError(stageName string, err error) OnStageComplete(stageName string, duration time.Duration) } // 内置观察者日志 Metrics type MetricsObserver struct { prometheus *prometheus.Registry } func (o *MetricsObserver) OnStageProgress(name string, processed, total int64) { stageProcessedGauge.WithLabelValues(name).Set(float64(processed)) }四、容错机制的四个层次层次一记录级容错。单条脏数据不影响整批处理脱敏后放入 Dead Letter Queue 供人工审核。层次二Stage 级容错。单个处理阶段失败时自动重试 3 次。3 次仍失败则跳过该阶段数据以不完整标记继续流转不阻塞下游。const maxStageRetries 3 func (p *Pipeline) runStageWithRetry( ctx context.Context, stage Stage, input -chan *Record, output chan- *Record, ) error { var lastErr error for attempt : 1; attempt maxStageRetries; attempt { // 每次重试需要新的 Channel retryInput : make(chan *Record, 1000) retryOutput : make(chan *Record, 1000) go func() { defer close(retryOutput) lastErr stage.Process(ctx, retryInput, retryOutput) }() // 将数据灌入重试管道 for rec : range input { retryInput - rec } close(retryInput) if lastErr nil { // 成功将结果传递给下游 for rec : range retryOutput { output - rec } return nil } p.logger.Warn(stage retry, slog.String(stage, stage.Name()), slog.Int(attempt, attempt), ) time.Sleep(time.Duration(attempt) * time.Second) } return fmt.Errorf(stage %s failed after %d retries: %w, stage.Name(), maxStageRetries, lastErr) }层次三管道级容错。整个管道异常时如配置错误告警通知开发者同时保留上下文快照便于排查。层次四基础设施级容错。数据库、外部 API 不可用时避免管道直接崩溃。通过连接池探测和健康检查提前发现问题。五、结语极简架构不是简陋架构。这个 ETL 管道虽然代码量不大但通过清晰的接口设计Stage、Filter、Strategy、Observer确保了每一层都有明确的职责边界和独立的容错能力。一年的运维数据侧面验证了这个设计的有效性月均故障 0 次单次 ETL 新增任务配置时间不超过 15 分钟新增一个数据源不需要重启任何服务。对于一个日均 200 万行的数据平台这些指标对应的是团队精力可以更多地投入到数据价值的挖掘而不是管道的维修。技术栈Go 1.22 / ClickHouse / Redis / MySQL / MongoDB

相关新闻