実験データのクリーニングパイプライン — 欠損値、重複、外れ値を再現性を持って処理
このトピックを終えると
教科書で学んだデータクリーニング、apply/map、mutable引数の落とし穴を組み合わせて、汚い実験データのCSVを再現可能な方法でクリーンにするパイプラインを自分で作成できます。毎回異なる方法でクリーニングして、結果が再現されないという問題がなくなります。
この記事は教育目的の一般的な例です。実務では、dbt、Great Expectations、Panderaなどの洗練されたツールを使用します。
「先週の結果はなぜ再現されないのか?」 — 手動クリーニングの落とし穴
毎週、実験データのCSVファイルを受け取り、分析するとしましょう。最初の週:
- Excelで開き、欠損セルを目視で確認
- 不自然に見える値を削除
- 複数のファイルから得られたデータを一つにまとめる
- 保存して分析
次の週も同じ手順を繰り返します。そして、指導教官が「先週の結果はなぜ今週再現されないのか?」と尋ねます。
原因は一つだけではありません。
問題1: 目視で削除した値が毎回異なります。ある週は3σ外の値を削除し、別の週は4σ外の値を削除します。
問題2: 欠損値の処理方法が毎回異なります。ある週は削除し、別の週は平均値で補完します。
問題3: どのファイルが今回欠けているのか、覚えていません。
正しいアプローチは、明示的なパイプラインです。各クリーニングステップをコードで明示的に記述し、そのコードを保存し、バージョン管理します。同じコードで同じデータを処理すると、正確に同じ結果が得られます。これが科学的な再現性の基本です。
ブラックボックスからコンポーネントへ
コンポーネント1:pandasの明示的なステップ
各クリーニングアクションを個別の関数として実装します。
import pandas as pdimport numpy as np
def load_and_validate(path: str, required_cols: list[str]) -> pd.DataFrame: df = pd.read_csv(path) missing = set(required_cols) - set(df.columns) if missing: raise ValueError(f"Missing columns: {missing}") return df
def remove_duplicates(df: pd.DataFrame, subset: list[str]) -> tuple[pd.DataFrame, int]: n_before = len(df) df_clean = df.drop_duplicates(subset=subset).copy() return df_clean, n_before - len(df_clean)
def drop_missing(df: pd.DataFrame, required_cols: list[str]) -> tuple[pd.DataFrame, int]: n_before = len(df) df_clean = df.dropna(subset=required_cols).copy() return df_clean, n_before - len(df_clean)
def filter_outliers_iqr(df: pd.DataFrame, col: str, factor: float = 1.5) -> tuple[pd.DataFrame, int]: q1 = df[col].quantile(0.25) q3 = df[col].quantile(0.75) iqr = q3 - q1 low, high = q1 - factor * iqr, q3 + factor * iqr n_before = len(df) df_clean = df[(df[col] >= low) & (df[col] <= high)].copy() return df_clean, n_before - len(df_clean)ポイント: 各関数が**どれだけ削除したか(カウント)**を返す。これを後でログとして記録します。
コンポーネント2:applyとmapによるカラム変換
値の変換は、applyまたはmapを使用して明示的に行います。
def normalize_gene_symbols(df: pd.DataFrame, col: str = "gene") -> pd.DataFrame: """ 遺伝子シンボルを大文字の標準形式に変換します。 """ df = df.copy() df[col] = df[col].str.upper().str.strip() return df
def standardize_units(df: pd.DataFrame, col: str, unit_col: str) -> pd.DataFrame: """ 濃度値と単位を組み合わせて、μMに統一します。 """ df = df.copy() def to_micromolar(row): value = row[col] unit = row[unit_col].lower() multiplier = {"nm": 0.001, "μm": 1, "um": 1, "mm": 1000, "m": 1_000_000} return value * multiplier.get(unit, np.nan) df["concentration_uM"] = df.apply(to_micromolar, axis=1) return df注意: 常にdf.copy()を使用します。元のDataFrameが変更されないようにします。
コンポーネント3:変更可能なトラップを回避
Python初心者が陥りやすい典型的な間違い:変更可能なデフォルト値。
# 危険なコードdef add_error_log(errors: list = []) -> list: errors.append("something wrong") return errors
# 最初の呼び出し: ["something wrong"]# 2回目の呼び出し: ["something wrong", "something wrong"] — 意図しない結果!理由:Python関数のデフォルト値は、関数が定義されたときに1回だけ評価されます。errors=[]の[]は、関数が定義されたときに作成され、その後も再利用されます。
修正:
def add_error_log(errors: list | None = None) -> list: if errors is None: errors = [] errors.append("something wrong") return errors同じ落とし穴はDataFrameにも存在します。
# 危険def process(df): df["new_col"] = df["old_col"] * 2 # 元のdfが変更される! return df
original = pd.read_csv("data.csv")processed = process(original)# original["new_col"] も存在する。元のデータが汚染される。
# 安全def process(df): df = df.copy() df["new_col"] = df["old_col"] * 2 return dfパイプラインの各ステップが純粋関数(入力を変更せずに新しい出力を返す)である場合にのみ、再現性が保たれます。
パイプラインの組み立て
ここで、上記のコンポーネントを1つのパイプラインにまとめます。
from dataclasses import dataclass, fieldfrom datetime import datetime
@dataclassclass CleaningReport: input_rows: int = 0 output_rows: int = 0 steps: list = field(default_factory=list) def add(self, step_name: str, removed: int, details: str = "") -> None: self.steps.append({ "step": step_name, "removed": removed, "details": details, "timestamp": datetime.now().isoformat() })
def clean_experiment_data( csv_path: str, required_cols: list[str] = None, outlier_col: str = "value", outlier_factor: float = 1.5) -> tuple[pd.DataFrame, CleaningReport]: required_cols = required_cols or ["gene", "value", "unit"] df = load_and_validate(csv_path, required_cols) report = CleaningReport(input_rows=len(df)) df = normalize_gene_symbols(df, "gene") report.add("normalize_gene_symbols", 0, "uppercased and stripped") df, n_dup = remove_duplicates(df, subset=["gene", "value"]) report.add("remove_duplicates", n_dup) df, n_missing = drop_missing(df, required_cols) report.add("drop_missing", n_missing) df = standardize_units(df, "value", "unit") report.add("standardize_units", 0) df, n_outlier = filter_outliers_iqr(df, outlier_col, outlier_factor) report.add("filter_outliers_iqr", n_outlier, f"factor={outlier_factor}") report.output_rows = len(df) return df, report使用例:
clean_df, report = clean_experiment_data("raw_data.csv")
print(f"入力: {report.input_rows}行, 出力: {report.output_rows}行")for step in report.steps: print(f" {step['step']}: -{step['removed']} ({step['details']})")期待される出力:
入力: 1000行, 出力: 934行
normalize_gene_symbols: -0 (uppercased and stripped)
remove_duplicates: -12
drop_missing: -30
standardize_units: -0
filter_outliers_iqr: -24 (factor=1.5)フェーディング — 埋めるべき3つの空白
空白1:スキーマ検証
各カラムの型と範囲を明示的に検証します。
def validate_schema(df: pd.DataFrame, schema: dict) -> list[str]: """ schema = {"col_name": {"type": float, "min": 0, "max": 100}} 検証に失敗した項目のエラーメッセージのリストを返します。空のリストであれば、すべて検証に合格です。 """ errors = [] for col, rules in schema.items(): if col not in df.columns: errors.append(f"Missing column: {col}") continue # TODO 1: 型の検証(df[col].dtype を確認) # TODO 2: min/max 範囲の検証(存在する場合、df[col].min()、df[col].max() を比較) pass return errorsヒント: if not pd.api.types.is_numeric_dtype(df[col]): errors.append(...). 範囲: if "min" in rules and df[col].min() < rules["min"]: errors.append(...).
空白2:ログファイルの保存
CleaningReport を CSV ログとして保存します。
def save_report(report: CleaningReport, log_path: str) -> None: """ report.steps を CSV として保存します。 """ # TODO: pd.DataFrame(report.steps) に変換した後、to_csv を使用 pass空白3:再現性ハッシュ
同一の入力に対して、同一のコードが同一の出力を生成することを確認するハッシュです。
import hashlib
def compute_output_hash(df: pd.DataFrame) -> str: """ df の内容をソートした後、ハッシュを計算します。 同じ df からは常に同じハッシュが生成される必要があります。 """ # TODO 1: df を安定した順序(例:すべてのカラムでソート)にソート # TODO 2: to_csv(index=False) で文字列に変換 # TODO 3: hashlib.sha256(bytes).hexdigest() を返す passヒント:
sorted_df = df.sort_values(list(df.columns)).reset_index(drop=True)content = sorted_df.to_csv(index=False).encode()return hashlib.sha256(content).hexdigest()[:16]考察 — 実際のデータパイプラインとの違い
dbt (data build tool):SQLベースのパイプラインの標準。各ステップはSQLモデルであり、依存関係が自動的に管理され、テストがパイプラインに統合されている。AirbnbやNetflixなどの企業で標準的に使用されている。
Great Expectations / Pandera:データスキーマと品質検証のためのフレームワーク。各カラムの期待値を宣言的に定義し、違反が発生した場合に通知する。
Apache Airflow / Prefect:大規模なパイプラインのオーケストレーション。スケジューリング、再試行、バックフィルなどの機能を提供する。
dask / polars:データが100GBを超える場合は、pandasの代わりにこれらのツールを使用する。並列処理と遅延評価をサポートする。
データリネージの追跡:実際の運用では、各出力がどの入力からどのような変換を経て生成されたかを自動的に追跡する。OpenLineageなどの標準が利用される。
拡張プロジェクト
1. Streamlitダッシュボード: クリーニングパイプラインをWeb UIで利用可能にする。ユーザーがCSVファイルをアップロードし、レポートを視覚化する。
2. Great Expectationsとの統合: CleaningReportの代わりに、Great Expectationsの標準レポートを使用する。
3. 複数ファイルの結合: 複数のCSVファイルを結合する際に、スキーマの一貫性を自動的に検証してから結合する。
4. データカタログ: 各実験データのメタデータをSQLiteに保存し、検索できるようにする。
このセクションの部品カタログ
- [F] データクリーニング: 欠損値、重複、外れ値を明示的に処理。各ステップのログを記録。
- [F] apply/map: 列変換の関数型アプローチ。axis=1 の理解。
- [F] 可変性の罠: 可変のデフォルト値、DataFrame のインプレース変更の回避。純粋な関数の原則。
- [W] ファイル I/O: CSV の読み書き(完成したスクリプトを提供)。
[F] = 自分で実装 / [W] = 完成したコードで提供。