パイプ&フィルタパターン — データの流れを設計する
工場の組み立てライン
「ユウキ、工場の組み立てラインって見たことある?」
「テレビで……部品が流れてきて、各ステーションで作業して、最後に完成品になる、みたいな?」
「それがパイプ&フィルタ」アヤカが言った。「データが流れてきて、各フィルターで変換・検証して、最後に完成品のデータになる」
具体的な例として、取引先からCSVで送られてくる注文データの取り込み処理を考えよう。
# 生CSV(汚い状態)
user_id,product_id,qty,note
" 101 ","P-001","two","急いで"
"102","P-002","3",""
"","P-003","1","通常配送"
"103","P-INVALID","5","特急"
このデータを最終的な Order オブジェクトに変換するまでに、複数のステップが必要だ。
Rack ミドルウェア:Rails の中のパイプ
Railsのリクエスト処理は、まさにパイプ&フィルタの実装。
# Rack ミドルウェアスタックを確認
rails middlewareuse 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 ApiRateLimiterMiddlewareRails でのパイプラインパターン
ビジネスロジックでもパイプ&フィルタを活用できる。
# 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
endAWS Step Functions でのパイプライン
より複雑な処理フローには AWS Step Functions を使う。
{
"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
endWARNING
フィルターチェーンが長くなると、どの段階でデータが変換されたかが追いにくくなる。各フィルターの前後でコンテキストの内容をログに記録し、デバッグ情報をコンテキストに含める設計にすると問題の特定が容易になる。
まとめ
| 適用場面 | 実装方法 |
|---|---|
| HTTPリクエスト処理 | Rackミドルウェア |
| データ変換・取り込み | カスタムパイプライン |
| 複雑なワークフロー | AWS Step Functions |
| バッチ処理 | Sidekiq チェーン / ActiveJob |
| マイクロサービス間のデータ流れ | AWS EventBridge Pipes |
「パイプ&フィルタって、知らないうちに使ってたんですね」ユウキが言った。「Rackミドルウェアがまさにそれだった」
「そう。パターンを知ると、既存のコードの意図が読めるようになる。それが一番の学習効果」
「フィルターが独立しているから、テストが書きやすかったです。各ステップを別々にテストできる」
「それが設計の良さが実感できる瞬間。テストしやすいコードは、変更しやすいコードでもある」
次章では、読み取りと書き込みを分離するCQRSパターンを学びます。パフォーマンス問題を解決する強力なアプローチです。