スループット向上のためのストリーミングパイプラインと名前付きパイプ
FIFOとプロセス置換を使い、中間ファイルなしでステージ間にデータをストリーミングします。
「スループット向上のためのストリーミングパイプラインと名前付きパイプ」はCoddyKit上の無料Linux Command Line & Bash Scripting Masteryレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはLinux Command Line & Bash Scripting Mastery学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Linux Command Line & Bash Scripting Masteryコースには全4レッスンが含まれています。
中間ファイルがスループットを低下させる理由
sort file.txt > tmp.txt && uniq tmp.txt > result.txtのようにコマンドを連結すると、ディスクへの書き込み、ディスクからの読み込み、そして次の段階が始まる前に最初の段階が完全に終了することによるパイプラインの停止という、隠れたコストが発生します。
ストリーミングパイプラインを使うと、このコストをなくせます。データは段階ごとに、メモリ上で生成側から消費側へ直接、並行して流れます。これがUnixパイプの基本的な考え方であり、名前付きパイプ(FIFO)はこの仕組みをさらに拡張します。
- 匿名パイプ(
|): 同じシェル行内で隣接する2つのコマンドを接続します。 - 名前付きパイプ(FIFO): 無関係なプロセス同士がストリーミングでデータをやり取りできる、ファイルシステム上の特殊なファイルです。
- プロセス置換: 別のコマンドの出力を、あるコマンドからファイルであるかのように扱えるようにします。
このレッスンでは、実際のBashワークフローでスループットを最大化するために、これら3つを適用する方法を説明します。
ストリーミングパイプラインの構造
匿名パイプは、あるプロセスの標準出力を次のプロセスの標準入力に接続します。カーネルは、通常Linuxでは64 KBの固定サイズのメモリ内バッファーを介して、両方のプロセスを同時に実行し続けます。
重要な点は、パイプラインの速度は最も遅い段階に左右されるということです。生成側の処理が速い場合は、バッファーが満杯になった時点でブロックします。消費側の処理が速い場合は、バッファーが空になった時点でブロックします。このバックプレッシャーは、追加コストなしで自動的に行われるフロー制御です。
以下の例では、一時ファイルを一度も書き込まずに、大きなアクセスログから一意なIPアドレスを数えます。各段階は並行して実行されます。
#!/usr/bin/env bash
# Stream a 2 GB access log — all stages run in parallel
grep '"GET' /var/log/nginx/access.log \
| awk '{print $1}' \
| sort \
| uniq -c \
| sort -rn \
| head -20mkfifoで名前付きパイプを作成する
名前付きパイプ(FIFO — First In, First Out)はmkfifoで作成します。通常のファイルと同じようにファイルシステム上に現れますが、書き込まれたデータがディスクに保存されることはなく、読み取り側のプロセスへ直接流れます。
次の動作を覚えておいてください。
- FIFOへの書き込みは、読み取り側がFIFOを開くまでブロックします。逆の場合も同様です。
- FIFOのエントリはファイルシステムに残るため、使い終わったら
rmで削除する必要があります。 - 複数の書き込み側を使用できますが、それらの間の順序は保証されません。
以下では、生成側がデータをFIFOに圧縮しながら、消費側が同時にS3へアップロードします。一時ファイルは必要ありません。
#!/usr/bin/env bash
mkfifo /tmp/stream_pipe
# Producer: compress in background
gzip -c /var/log/syslog > /tmp/stream_pipe &
# Consumer: read from FIFO (runs in foreground)
wc -l < /tmp/stream_pipe
wait
rm /tmp/stream_pipeプロセス置換: コマンドをファイルとして扱う
プロセス置換では、<(command)または>(command)という構文を使用します。Bashは内部でFIFO(または/dev/fd/Nファイルディスクリプター)を作成し、そのパスを外側のコマンドに渡します。
これは、標準入力ではなくファイル名引数を必要とするツールで特に便利です。プロセス置換を使わない場合は一時ファイルが必要ですが、使えば直接ストリーミングできます。
<(cmd)— 外側のコマンドがcmdの出力を読み取ります。>(cmd)— 外側のコマンドがcmdの入力へ書き込みます。
#!/usr/bin/env bash
# diff two sorted streams without creating temp files
diff <(sort /etc/passwd) <(sort /etc/group)
# Compare live command output against a baseline
diff <(ls /usr/bin | sort) <(cat ~/bin_baseline.txt | sort)tee: ストリームを複数の消費側に分岐する
teeは標準入力を読み取り、標準出力と1つ以上のファイルの両方に書き込みます。プロセス置換と組み合わせると、1つのストリームを複数の処理パイプラインに同時に分岐できます。ディスクに触れる必要もありません。
たとえば、生データをログに記録しながら同時に処理したい場合に、このパターンが役立ちます。
#!/usr/bin/env bash
# Generate 100000 random numbers, then simultaneously:
# 1. compute the sum
# 2. find the maximum
# 3. count lines (saved to a variable)
seq 1 100000 \
| tee >(awk '{s+=$1} END{print "Sum:", s}') \
>(awk 'BEGIN{m=0} $1>m{m=$1} END{print "Max:", m}') \
| wc -l | xargs echo "Count:"ファンアウトパターン: 1つの生成側、多数の消費側
1つのデータソースから複数の独立した消費側へデータを供給する必要がある場合は、teeと複数の>()プロセス置換を組み合わせます。各消費側はストリーム全体を受け取り、並行して実行されます。
これにより、ソースファイルを複数回読み取る必要がなくなります。10 GBのファイルではその差は非常に大きく、N回のディスク読み取りではなく1回で済みます。
#!/usr/bin/env bash
# Read a large CSV once; simultaneously:
# - count rows
# - extract column 2 to a file
# - pass column 3 to a stats script
cat large_data.csv \
| tee \
>(wc -l > /tmp/row_count.txt) \
>(cut -d',' -f2 > /tmp/col2.txt) \
>(cut -d',' -f3 | awk '{sum+=$1} END{print sum}' > /tmp/col3_sum.txt) \
> /dev/null
echo "Rows:" $(cat /tmp/row_count.txt)
echo "Col3 sum:" $(cat /tmp/col3_sum.txt)ファンインパターン: 多数の生成側、1つの消費側
ファンアウトの逆がファンインです。複数の独立したソースから、1つの消費側へストリーミングします。名前付きFIFOを使うと簡単に実現できます。
一般的な用途として、複数のサーバーからのログストリームをリアルタイムで統合したり、並列ワーカーからの部分的な結果を集約したりする場合があります。
書き込み側が複数ある場合、消費側からは出力がインターリーブして見えます。各行が独立している行指向のデータでは問題ありませんが、順序が重要な場合は自分で順序を処理する必要があります。
#!/usr/bin/env bash
mkfifo /tmp/fanin_pipe
# Three producers write concurrently into the same FIFO
for host in web1 web2 web3; do
ssh "$host" 'tail -n 500 /var/log/app.log' > /tmp/fanin_pipe &
done
# Single consumer reads all merged output
grep 'ERROR' /tmp/fanin_pipe | sort | uniq -c | sort -rn
wait
rm /tmp/fanin_pipemkfifoを使った並列圧縮
FIFOの最も実用的な用途の1つが並列圧縮です。pigz(並列gzip)やpbzip2などのツールはストリームを読み取るため、非圧縮ファイルを一時保存せずに生データを直接パイプできます。
以下のパターンでは、ディレクトリをアーカイブし、すべてのCPUコアで圧縮し、その結果をリモートホストへストリーミングします。これらはすべて同時に実行されます。
#!/usr/bin/env bash
# Tar + parallel compress + stream to remote — no temp files
# Requires: pigz (parallel gzip)
tar cf - /data/large_dir \
| pigz -p 4 \
| ssh backup-host 'cat > /backups/large_dir.tar.gz'
# Verify the remote file exists
ssh backup-host 'ls -lh /backups/large_dir.tar.gz'バッファーサイズとブロッキングの制御
パイプにはカーネルバッファー(通常64 KB)があります。バッファーが満杯になると書き込み側がブロックし、空になると読み取り側がブロックします。通常は望ましい動作ですが、場合によってはブロッキングによってデッドロックが発生します。
デッドロックの危険: プロセスAがFIFO1に書き込み、FIFO2から読み取る一方で、プロセスBがFIFO2に書き込み、FIFO1から読み取る場合、互いが先にデータを消費するのを待って、両方がブロックする可能性があります。
解決策:
- 少なくとも一方をバックグラウンド(
&)で実行し、シェルをブロックしないようにします。 mbufferまたはpvを使い、段階の間により大きなメモリ内バッファーを追加します。pv -q -B 128mを使って128 MBのバッファーを挿入し、スループットの急激な変動を平滑化します。
#!/usr/bin/env bash
# pv adds a 64 MB buffer and shows throughput
# Useful when producer and consumer have bursty speeds
dd if=/dev/urandom bs=1M count=200 \
| pv -B 64m \
| gzip \
| wc -c実践例: リアルタイムログアグリゲーター
完全で現実的なパターンを見てみましょう。複数のログファイルをtailし、名前付きパイプを通してストリームを統合し、エラーをフィルタリングして、ライブサマリーを書き込みます。中間ファイルは一切使わず、すべての段階を並列に実行します。
これは、本番サーバーでバックグラウンドの監視スクリプトとして実行するようなパイプラインです。
#!/usr/bin/env bash
FIFO=/tmp/log_aggregator
mkfifo "$FIFO"
cleanup() { rm -f "$FIFO"; }
trap cleanup EXIT INT TERM
# Fan-in: tail multiple logs into the FIFO
tail -F /var/log/syslog /var/log/auth.log > "$FIFO" &
TAIL_PID=$!
# Consumer: filter and timestamp errors in real time
grep --line-buffered -i 'error\|fail\|crit' "$FIFO" \
| while IFS= read -r line; do
printf '[%s] %s\n' "$(date '+%H:%M:%S')" "$line"
done
kill "$TAIL_PID" 2>/dev/nullパイプラインと一時ファイルのベンチマーク
timeを使うと、ストリーミング方式と一時ファイル方式の実際のスループットの違いを測定できます。大規模なデータセットでは、次の理由からパイプライン方式が優れています。
- 段階が並行して実行されるため、CPU処理とI/Oが重なります。
- 中間データのディスクI/Oがなく、最終出力だけがディスクに書き込まれます。
- ストリームとして処理するため、入力サイズに関係なくメモリ使用量が一定です。
両方の方式を比較する簡単なベンチマークは次のとおりです。
#!/usr/bin/env bash
# Approach 1: Temp file (sequential)
time bash -c '
seq 1 5000000 > /tmp/nums.txt
sort -n /tmp/nums.txt > /tmp/sorted.txt
uniq /tmp/sorted.txt | wc -l
rm /tmp/nums.txt /tmp/sorted.txt
'
echo '---'
# Approach 2: Streaming pipeline (concurrent)
time bash -c 'seq 1 5000000 | sort -n | uniq | wc -l'理解度チェック: 名前付きパイプのブロッキング動作
次のスクリプトについて考えてみましょう。
mkfifo /tmp/mypipe
echo 'hello' > /tmp/mypipe
echo 'done'バックグラウンドプロセスも読み取り側もない状態でこのスクリプトを実行すると、何が起こるでしょうか。
振り返り:ストリーミングパイプラインと名前付きパイプ
このレッスンでは、中間ファイルを使わずにプロセス間でデータを効率よく移動する方法を学びました。
- 匿名パイプ(
|)は隣接するコマンドを接続し、バックプレッシャーを自動的に適用しながら、すべての段階を同時に実行します。 - 名前付きパイプ(
mkfifo)はFIFOのファイルシステムエントリを作成し、無関係なプロセスやバックグラウンドプロセス同士でストリームをやり取りできるようにします。読み取り側が存在するまで書き込みはブロックされます。 - プロセス置換(
<(cmd)、>(cmd))を使うと、ファイル名を受け取るコマンドでストリームを透過的に消費または生成できます。 tee+>()を使うと、ソースを再読み込みせずに、1つのストリームを複数の並行コンシューマーへ分配できます。- Fan-inでは、共有FIFOを介して複数のプロデューサーからの出力を1つのコンシューマーに統合します。
- バックプレッシャーとブロッキングはバグではなく機能です。ただし、FIFOのペアでデッドロックが発生しないよう、必ず少なくとも片方をバックグラウンドで実行してください。
- 各段階の処理量に波がある場合は、
pvまたはmbufferを使ってバッファーを大きくし、スループットを監視してください。
これらの手法は、高スループットなBashデータエンジニアリングの基盤です。一定のメモリ使用量でGB単位のデータを処理し、CPUとI/Oの並列性を最大限に高められます。
よくある質問
「スループット向上のためのストリーミングパイプラインと名前付きパイプ」レッスンは無料ですか?
はい。「スループット向上のためのストリーミングパイプラインと名前付きパイプ」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Linux Command Line & Bash Scripting Masteryコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Linux Command Line & Bash Scripting Masteryコースには全4レッスンが含まれています。
「スループット向上のためのストリーミングパイプラインと名前付きパイプ」で何を学びますか?
FIFOとプロセス置換を使い、中間ファイルなしでステージ間にデータをストリーミングします。 ブラウザで直接実行するハンズオンコードでLinux Command Line & Bash Scripting Masteryを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Linux Command Line & Bash Scripting Masteryを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのLinux Command Line & Bash Scripting Masteryは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。
「スループット向上のためのストリーミングパイプラインと名前付きパイプ」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このLinux Command Line & Bash Scripting Masteryレッスンでコードを書いて実行できますか?
はい。すべてのLinux Command Line & Bash Scripting Masteryレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- スクリプトのプロファイリングと不要なサブシェルの回避
- xargs -Pとバックグラウンドジョブによる並列処理
- GNU parallelによるワークロードのオーケストレーション
- スループット向上のためのストリーミングパイプラインと名前付きパイプ