mybook

パイプ&フィルタパターン — データの流れを設計する

工場の組み立てライン

「ユウキ、工場の組み立てラインって見たことある?」

「テレビで……部品が流れてきて、各ステーションで作業して、最後に完成品になる、みたいな?」

「それがパイプ&フィルタ」アヤカが言った。「データが流れてきて、各フィルターで変換・検証して、最後に完成品のデータになる」

具体的な例として、取引先からCSVで送られてくる注文データの取り込み処理を考えよう。

# 生CSV(汚い状態)
user_id,product_id,qty,note
" 101 ","P-001","two","急いで"
"102","P-002","3",""
"","P-003","1","通常配送"
"103","P-INVALID","5","特急"

このデータを最終的な Order オブジェクトに変換するまでに、複数のステップが必要だ。

Loading diagram...

Rack ミドルウェア:Rails の中のパイプ

Railsのリクエスト処理は、まさにパイプ&フィルタの実装。

# Rack ミドルウェアスタックを確認
rails middleware
use ActionDispatch::HostAuthorization
use Rack::Sendfile
use ActionDispatch::Static
use ActionDispatch::Executor
use ActiveSupport::Cache::Strategy::LocalCache::Middleware
use Rack::Runtime
use ActionDispatch::RequestId
use ActionDispatch::RemoteIp
use Rails::Rack::Logger
use ActionDispatch::ShowExceptions
use ActionDispatch::DebugExceptions
use ActionDispatch::ActionableExceptions
use ActionDispatch::Reloader
use ActionDispatch::Callbacks
use ActiveRecord::Migration::CheckPending
use ActionDispatch::Cookies
use ActionDispatch::Session::CookieStore
use ActionDispatch::Flash
use ActionDispatch::ContentSecurityPolicy::Middleware
use Rack::Head
use Rack::ConditionalGet
use Rack::ETag
use Rack::TempfileReaper
run MyApp::Application.routes

リクエストはこのスタックを上から順に通り、レスポンスは逆順に通る。

INFO

各ミドルウェアは initialize(app)call(env) の2メソッドを持つシンプルなインターフェース。フィルターはこの契約を守るだけでよく、他のフィルターの存在を知らない。追加も削除も1行で済む。

カスタムRackミドルウェアの実装

# lib/middleware/request_logger_middleware.rb
class RequestLoggerMiddleware
  def initialize(app)
    @app = app
  end
 
  def call(env)
    start_time = Time.current
    request = Rack::Request.new(env)
 
    # リクエストログ(入力フィルタリング)
    request_id = env["action_dispatch.request_id"] || SecureRandom.uuid
    Rails.logger.info({
      type: "request_start",
      request_id: request_id,
      method: request.request_method,
      path: request.path,
      remote_ip: request.ip
    }.to_json)
 
    # 次のミドルウェアへ(パイプ)
    status, headers, body = @app.call(env)
 
    # レスポンスログ(出力フィルタリング)
    duration_ms = ((Time.current - start_time) * 1000).round(2)
    Rails.logger.info({
      type: "request_end",
      request_id: request_id,
      status: status,
      duration_ms: duration_ms
    }.to_json)
 
    [status, headers, body]
  end
end
# lib/middleware/api_authentication_middleware.rb
class ApiAuthenticationMiddleware
  SKIP_PATHS = %w[/health /api/v1/sessions /api/v1/users].freeze
 
  def initialize(app)
    @app = app
  end
 
  def call(env)
    request = Rack::Request.new(env)
 
    # 認証不要なパスはスキップ
    return @app.call(env) if skip_auth?(request)
    return @app.call(env) unless api_request?(request)
 
    token = extract_token(env)
 
    unless valid_token?(token)
      return [
        401,
        { "Content-Type" => "application/json" },
        [{ error: "Unauthorized", message: "Invalid or missing API token" }.to_json]
      ]
    end
 
    # 認証済みユーザー情報を後続のミドルウェアに渡す
    env["current_user_id"] = decode_token(token)[:user_id]
    @app.call(env)
  end
 
  private
 
  def skip_auth?(request)
    SKIP_PATHS.any? { |path| request.path.start_with?(path) }
  end
 
  def api_request?(request)
    request.path.start_with?("/api/")
  end
 
  def extract_token(env)
    auth_header = env["HTTP_AUTHORIZATION"]
    auth_header&.split(" ")&.last
  end
 
  def valid_token?(token)
    token.present? && JsonWebToken.decode(token)
  rescue JWT::DecodeError
    false
  end
 
  def decode_token(token)
    JsonWebToken.decode(token)
  end
end
# lib/middleware/api_rate_limiter_middleware.rb
class ApiRateLimiterMiddleware
  LIMIT_PER_MINUTE = 60
  BURST_LIMIT = 100
 
  def initialize(app)
    @app = app
  end
 
  def call(env)
    request = Rack::Request.new(env)
    return @app.call(env) unless api_request?(request)
 
    client_id = env["current_user_id"] || env["REMOTE_ADDR"]
    window = Time.current.to_i / 60
    key = "rate_limit:#{client_id}:#{window}"
 
    count = Rails.cache.increment(key, 1, expires_in: 90.seconds, raw: true)
 
    remaining = [LIMIT_PER_MINUTE - count, 0].max
    reset_at = (window + 1) * 60
 
    if count > BURST_LIMIT
      return [
        429,
        {
          "Content-Type" => "application/json",
          "X-RateLimit-Limit" => LIMIT_PER_MINUTE.to_s,
          "X-RateLimit-Remaining" => "0",
          "X-RateLimit-Reset" => reset_at.to_s,
          "Retry-After" => (reset_at - Time.current.to_i).to_s
        },
        [{ error: "Rate limit exceeded", retry_after: reset_at }.to_json]
      ]
    end
 
    status, headers, body = @app.call(env)
 
    headers["X-RateLimit-Limit"] = LIMIT_PER_MINUTE.to_s
    headers["X-RateLimit-Remaining"] = remaining.to_s
    headers["X-RateLimit-Reset"] = reset_at.to_s
 
    [status, headers, body]
  end
 
  private
 
  def api_request?(request)
    request.path.start_with?("/api/")
  end
end
 
# config/application.rb に登録
# config.middleware.insert_before ActionDispatch::RequestId, RequestLoggerMiddleware
# config.middleware.use ApiAuthenticationMiddleware
# config.middleware.use ApiRateLimiterMiddleware

Rails でのパイプラインパターン

ビジネスロジックでもパイプ&フィルタを活用できる。

# app/pipelines/pipeline.rb
# 汎用パイプライン基底クラス
class Pipeline
  def initialize
    @filters = []
  end
 
  def use(filter)
    @filters << filter
    self
  end
 
  def call(context)
    @filters.reduce(context) do |ctx, filter|
      break ctx if ctx[:halted]
      filter.call(ctx)
    end
  end
end
# app/pipelines/order_import_pipeline.rb
# CSVファイルからの注文取り込みパイプライン
class OrderImportPipeline
  Result = Data.define(:success?, :processed_count, :skipped_count, :errors)
 
  FILTERS = [
    Filters::CsvParseFilter,
    Filters::ColumnValidationFilter,
    Filters::DataNormalizationFilter,
    Filters::UserValidationFilter,
    Filters::ProductValidationFilter,
    Filters::DuplicateCheckFilter,
    Filters::OrderCreationFilter
  ].freeze
 
  def self.process(csv_content)
    new(csv_content).process
  end
 
  def initialize(csv_content)
    @csv_content = csv_content
  end
 
  def process
    context = {
      raw_data: @csv_content,
      rows: [],
      normalized_rows: [],
      valid_rows: [],
      orders: [],
      errors: [],
      skipped: 0,
      halted: false
    }
 
    final_context = FILTERS.reduce(context) do |ctx, filter_class|
      break ctx if ctx[:halted]
      filter_class.new.call(ctx)
    end
 
    Result.new(
      success?: final_context[:errors].empty?,
      processed_count: final_context[:orders].size,
      skipped_count: final_context[:skipped],
      errors: final_context[:errors]
    )
  end
end
# app/pipelines/filters/csv_parse_filter.rb
module Filters
  class CsvParseFilter
    def call(context)
      rows = CSV.parse(
        context[:raw_data],
        headers: true,
        skip_blanks: true,
        liberal_parsing: true
      ).map(&:to_h)
 
      context.merge(rows: rows)
    rescue CSV::MalformedCSVError => e
      context.merge(
        errors: context[:errors] + ["CSV形式エラー: #{e.message}"],
        halted: true
      )
    end
  end
end
# app/pipelines/filters/column_validation_filter.rb
module Filters
  class ColumnValidationFilter
    REQUIRED_COLUMNS = %w[user_id product_id qty].freeze
 
    def call(context)
      return context if context[:rows].empty?
 
      actual_columns = context[:rows].first&.keys || []
      missing = REQUIRED_COLUMNS - actual_columns
 
      if missing.any?
        return context.merge(
          errors: context[:errors] + ["必須列が不足しています: #{missing.join(', ')}"],
          halted: true
        )
      end
 
      context
    end
  end
end
# app/pipelines/filters/data_normalization_filter.rb
module Filters
  class DataNormalizationFilter
    QTY_WORDS = { "one" => 1, "two" => 2, "three" => 3 }.freeze
 
    def call(context)
      normalized = context[:rows].map.with_index(1) do |row, line_no|
        normalize_row(row, line_no)
      end
 
      valid = normalized.select { |r| r[:valid] }
      invalid = normalized.reject { |r| r[:valid] }
 
      new_errors = invalid.map { |r| "行#{r[:line_no]}: #{r[:error]}" }
 
      context.merge(
        normalized_rows: valid,
        errors: context[:errors] + new_errors,
        skipped: context[:skipped] + invalid.size
      )
    end
 
    private
 
    def normalize_row(row, line_no)
      user_id = row["user_id"]&.strip&.to_i
      product_id = row["product_id"]&.strip
      qty_raw = row["qty"]&.strip&.downcase
 
      return { valid: false, line_no: line_no, error: "user_id が空です" } if user_id.zero?
      return { valid: false, line_no: line_no, error: "product_id が空です" } if product_id.blank?
 
      qty = QTY_WORDS[qty_raw] || qty_raw.to_i
 
      return { valid: false, line_no: line_no, error: "qty が無効な値です: #{qty_raw}" } if qty <= 0
 
      { valid: true, line_no: line_no, user_id: user_id, product_id: product_id, qty: qty, note: row["note"]&.strip }
    end
  end
end
# app/pipelines/filters/duplicate_check_filter.rb
module Filters
  class DuplicateCheckFilter
    def call(context)
      rows = context[:normalized_rows]
 
      # 今日の重複注文チェック
      existing_keys = Set.new(
        Order.where(created_at: Time.current.beginning_of_day..)
             .where(user_id: rows.map { |r| r[:user_id] }.uniq)
             .select(:user_id, :import_key)
             .map { |o| o.import_key }
      )
 
      unique_rows = rows.reject do |row|
        key = import_key(row)
        existing_keys.include?(key)
      end
 
      skipped = rows.size - unique_rows.size
 
      context.merge(
        normalized_rows: unique_rows,
        skipped: context[:skipped] + skipped
      )
    end
 
    private
 
    def import_key(row)
      "#{row[:user_id]}-#{row[:product_id]}-#{Date.today}"
    end
  end
end
# app/pipelines/filters/order_creation_filter.rb
module Filters
  class OrderCreationFilter
    def call(context)
      created_orders = []
      row_errors = []
 
      context[:normalized_rows].each do |row|
        begin
          result = OrderCreationService.new(
            user: User.find(row[:user_id]),
            items: [{ product_id: Product.find_by(sku: row[:product_id]).id, quantity: row[:qty] }]
          ).call
 
          if result.success?
            created_orders << result.order
          else
            row_errors << "行#{row[:line_no]}: #{result.errors.join(', ')}"
          end
        rescue => e
          row_errors << "行#{row[:line_no]}: #{e.message}"
        end
      end
 
      context.merge(
        orders: created_orders,
        errors: context[:errors] + row_errors
      )
    end
  end
end
# 使用例
result = OrderImportPipeline.process(params[:file].read)
 
if result.success?
  render json: {
    message: "取り込み完了",
    processed: result.processed_count,
    skipped: result.skipped_count
  }
else
  render json: {
    message: "取り込み部分完了",
    processed: result.processed_count,
    skipped: result.skipped_count,
    errors: result.errors
  }, status: :unprocessable_entity
end

AWS Step Functions でのパイプライン

より複雑な処理フローには AWS Step Functions を使う。

Loading diagram...
{
  "Comment": "注文インポートパイプライン",
  "StartAt": "ValidateInput",
  "States": {
    "ValidateInput": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:validate-input",
      "ResultPath": "$.validation_result",
      "Next": "NormalizeData",
      "Catch": [
        {
          "ErrorEquals": ["ValidationError"],
          "ResultPath": "$.error",
          "Next": "HandleError"
        }
      ],
      "Retry": [
        {
          "ErrorEquals": ["Lambda.ServiceException"],
          "IntervalSeconds": 2,
          "MaxAttempts": 3,
          "BackoffRate": 2
        }
      ]
    },
    "NormalizeData": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:normalize-data",
      "ResultPath": "$.normalized_data",
      "Next": "ProcessOrder"
    },
    "ProcessOrder": {
      "Type": "Task",
      "Resource": "arn:aws:states:::ecs:runTask.sync",
      "Parameters": {
        "LaunchType": "FARGATE",
        "Cluster": "arn:aws:ecs:ap-northeast-1:xxx:cluster/myapp",
        "TaskDefinition": "arn:aws:ecs:ap-northeast-1:xxx:task-definition/order-processor",
        "Overrides": {
          "ContainerOverrides": [
            {
              "Name": "app",
              "Environment": [
                {
                  "Name": "ORDER_DATA",
                  "Value.$": "States.JsonToString($.normalized_data)"
                }
              ]
            }
          ]
        }
      },
      "Next": "SendNotification",
      "Catch": [
        {
          "ErrorEquals": ["States.TaskFailed"],
          "Next": "HandleError"
        }
      ]
    },
    "SendNotification": {
      "Type": "Task",
      "Resource": "arn:aws:states:::sns:publish",
      "Parameters": {
        "TopicArn": "arn:aws:sns:ap-northeast-1:xxx:order-events",
        "Message.$": "States.Format('注文インポート完了: {}件', $.processed_count)"
      },
      "End": true
    },
    "HandleError": {
      "Type": "Task",
      "Resource": "arn:aws:lambda:ap-northeast-1:xxx:function:handle-error",
      "End": true
    }
  }
}

Golang でのパイプラインパターン

Golangでは、関数型のパイプラインが書きやすい。

// pipeline/pipeline.go
package pipeline
 
import "context"
 
type Context map[string]interface{}
 
type Filter func(ctx context.Context, data Context) (Context, error)
 
type Pipeline struct {
    filters []Filter
}
 
func New(filters ...Filter) *Pipeline {
    return &Pipeline{filters: filters}
}
 
func (p *Pipeline) Execute(ctx context.Context, initial Context) (Context, error) {
    current := initial
 
    for _, filter := range p.filters {
        result, err := filter(ctx, current)
        if err != nil {
            return current, err
        }
        current = result
    }
 
    return current, nil
}
// pipeline/filters.go
package pipeline
 
import (
    "context"
    "encoding/csv"
    "fmt"
    "strconv"
    "strings"
)
 
func ParseCSVFilter(ctx context.Context, data Context) (Context, error) {
    rawData, ok := data["raw_data"].(string)
    if !ok {
        return data, fmt.Errorf("raw_data が見つかりません")
    }
 
    reader := csv.NewReader(strings.NewReader(rawData))
    reader.TrimLeadingSpace = true
    records, err := reader.ReadAll()
    if err != nil {
        return data, fmt.Errorf("CSV解析エラー: %w", err)
    }
 
    data["records"] = records
    return data, nil
}
 
func ValidateColumnsFilter(ctx context.Context, data Context) (Context, error) {
    records, _ := data["records"].([][]string)
    if len(records) == 0 {
        return data, fmt.Errorf("CSVが空です")
    }
 
    headers := records[0]
    required := []string{"user_id", "product_id", "qty"}
 
    for _, req := range required {
        found := false
        for _, h := range headers {
            if h == req {
                found = true
                break
            }
        }
        if !found {
            return data, fmt.Errorf("必須列が不足: %s", req)
        }
    }
 
    return data, nil
}
 
func NormalizeDataFilter(ctx context.Context, data Context) (Context, error) {
    records, _ := data["records"].([][]string)
    headers := records[0]
 
    type Row struct {
        UserID    int64
        ProductID string
        Qty       int
    }
 
    var normalized []Row
    var errors []string
 
    for i, record := range records[1:] {
        row := make(map[string]string)
        for j, h := range headers {
            if j < len(record) {
                row[h] = strings.TrimSpace(record[j])
            }
        }
 
        userID, err := strconv.ParseInt(row["user_id"], 10, 64)
        if err != nil || userID == 0 {
            errors = append(errors, fmt.Sprintf("行%d: user_id が無効", i+2))
            continue
        }
 
        qty, err := strconv.Atoi(row["qty"])
        if err != nil || qty <= 0 {
            errors = append(errors, fmt.Sprintf("行%d: qty が無効", i+2))
            continue
        }
 
        normalized = append(normalized, Row{
            UserID:    userID,
            ProductID: row["product_id"],
            Qty:       qty,
        })
    }
 
    data["normalized"] = normalized
    data["errors"] = errors
    return data, nil
}

パイプラインのテスト

フィルターが独立しているため、単体テストが書きやすい。

# spec/pipelines/filters/data_normalization_filter_spec.rb
RSpec.describe Filters::DataNormalizationFilter do
  subject(:filter) { described_class.new }
 
  context "正常なデータ" do
    let(:context) do
      {
        rows: [
          { "user_id" => " 101 ", "product_id" => "P-001", "qty" => "2", "note" => " 急ぎ " }
        ],
        errors: [],
        skipped: 0
      }
    end
 
    it "空白が除去されて数値に変換される" do
      result = filter.call(context)
      expect(result[:normalized_rows].first[:user_id]).to eq(101)
      expect(result[:normalized_rows].first[:qty]).to eq(2)
      expect(result[:normalized_rows].first[:note]).to eq("急ぎ")
    end
  end
 
  context "qtyが英語表記の場合" do
    let(:context) do
      {
        rows: [{ "user_id" => "101", "product_id" => "P-001", "qty" => "two" }],
        errors: [],
        skipped: 0
      }
    end
 
    it "数値に変換される" do
      result = filter.call(context)
      expect(result[:normalized_rows].first[:qty]).to eq(2)
    end
  end
 
  context "無効なqtyの場合" do
    let(:context) do
      {
        rows: [{ "user_id" => "101", "product_id" => "P-001", "qty" => "invalid" }],
        errors: [],
        skipped: 0
      }
    end
 
    it "スキップされてエラーが追加される" do
      result = filter.call(context)
      expect(result[:normalized_rows]).to be_empty
      expect(result[:errors]).not_to be_empty
      expect(result[:skipped]).to eq(1)
    end
  end
end
# spec/pipelines/order_import_pipeline_spec.rb
RSpec.describe OrderImportPipeline do
  let(:valid_csv) do
    <<~CSV
      user_id,product_id,qty,note
      101,P-001,2,通常配送
      102,P-002,1,
    CSV
  end
 
  before do
    create(:user, id: 101)
    create(:user, id: 102)
    create(:product, sku: "P-001", stock: 10, price: 1000)
    create(:product, sku: "P-002", stock: 5, price: 500)
  end
 
  it "有効なCSVから注文が作成される" do
    result = described_class.process(valid_csv)
 
    expect(result.success?).to be true
    expect(result.processed_count).to eq(2)
  end
 
  context "不正なCSVの場合" do
    let(:invalid_csv) { "not,a,valid\ncsv" * 100 }
 
    it "エラーを返す" do
      result = described_class.process("user_id\n101")
 
      expect(result.success?).to be false
    end
  end
end

WARNING

フィルターチェーンが長くなると、どの段階でデータが変換されたかが追いにくくなる。各フィルターの前後でコンテキストの内容をログに記録し、デバッグ情報をコンテキストに含める設計にすると問題の特定が容易になる。

まとめ

適用場面実装方法
HTTPリクエスト処理Rackミドルウェア
データ変換・取り込みカスタムパイプライン
複雑なワークフローAWS Step Functions
バッチ処理Sidekiq チェーン / ActiveJob
マイクロサービス間のデータ流れAWS EventBridge Pipes

「パイプ&フィルタって、知らないうちに使ってたんですね」ユウキが言った。「Rackミドルウェアがまさにそれだった」

「そう。パターンを知ると、既存のコードの意図が読めるようになる。それが一番の学習効果」

「フィルターが独立しているから、テストが書きやすかったです。各ステップを別々にテストできる」

「それが設計の良さが実感できる瞬間。テストしやすいコードは、変更しやすいコードでもある」


次章では、読み取りと書き込みを分離するCQRSパターンを学びます。パフォーマンス問題を解決する強力なアプローチです。