Конвейер данных последовательно передаёт результат одного этапа следующему и помогает разделить сложную обработку на небольшие операции.
Что такое Data Pipeline?
Data Pipeline (Конвейер данных) — последовательная цепочка функций, где выход одной становится входом другой.
numbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
# Без pipeline: промежуточные переменные
even = [x for x in numbers if x % 2 == 0]
doubled = [x * 2 for x in even]
total = sum(doubled)
print(total) # 60
# С pipeline: одна цепочка
from functools import reduce
total = reduce(
lambda acc, x: acc + x,
map(lambda x: x * 2, filter(lambda x: x % 2 == 0, numbers)),
0
)
print(total) # 60
Визуализация
[1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
↓ filter(чётные)
[2, 4, 6, 8, 10]
↓ map(удвоить)
[4, 8, 12, 16, 20]
↓ reduce(сумма)
60
Базовый Pipeline: filter → map → reduce
from functools import reduce
numbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
# Сумма квадратов чётных чисел
result = reduce(
lambda acc, x: acc + x, # Шаг 3: сумма
map(
lambda x: x ** 2, # Шаг 2: квадрат
filter(lambda x: x % 2 == 0, numbers) # Шаг 1: чётные
),
0
)
print(result) # 220
# Чётные: [2, 4, 6, 8, 10]
# Квадраты: [4, 16, 36, 64, 100]
# Сумма: 220
Тот же pipeline читается удобнее через list comprehension:
result = sum([x ** 2 for x in numbers if x % 2 == 0])
print(result) # 220
Функция pipe()
Универсальная функция для построения pipeline:
def pipe(data, *functions):
"""Применить функции последовательно."""
result = data
for func in functions:
result = func(result)
return result
# Определяем шаги
def filter_even(numbers):
return [x for x in numbers if x % 2 == 0]
def double_all(numbers):
return [x * 2 for x in numbers]
def sum_all(numbers):
return sum(numbers)
# Запускаем pipeline
numbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
result = pipe(numbers, filter_even, double_all, sum_all)
print(result) # 60
Каждый шаг — отдельная функция. Легко добавлять, убирать, переставлять шаги.
Практический пример: обработка заказов
from functools import reduce
data = [
{"product": "Phone", "quantity": 2, "price": 500},
{"product": "Laptop", "quantity": 1, "price": 1200},
{"product": "Mouse", "quantity": 5, "price": 25},
{"product": "Monitor", "quantity": 2, "price": 300}
]
# Pipeline: quantity > 1 → вычислить total → суммировать
filtered = filter(lambda item: item["quantity"] > 1, data)
totals = map(lambda item: item["quantity"] * item["price"], filtered)
grand_total = reduce(lambda acc, x: acc + x, totals, 0)
print(grand_total) # 1725
# Phone: 2 × 500 = 1000
# Mouse: 5 × 25 = 125
# Monitor: 2 × 300 = 600
# Итого: 1725
Класс Pipeline
Fluent interface для удобной работы:
class Pipeline:
def __init__(self, data):
self.data = data
def filter(self, predicate):
self.data = [x for x in self.data if predicate(x)]
return self
def map(self, transform):
self.data = [transform(x) for x in self.data]
return self
def reduce(self, reducer, initial=None):
from functools import reduce
if initial is None:
return reduce(reducer, self.data)
return reduce(reducer, self.data, initial)
def collect(self):
return self.data
# Использование
numbers = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]
result = (Pipeline(numbers)
.filter(lambda x: x % 2 == 0)
.map(lambda x: x * 2)
.reduce(lambda acc, x: acc + x, 0))
print(result) # 60
Частые ошибки
Ошибка 1: Изменение оригинальных данных
# ❌ Плохо: sort() меняет исходный список
def bad_pipeline(data):
data.sort()
return [x * 2 for x in data]
numbers = [3, 1, 2]
result = bad_pipeline(numbers)
print(numbers) # [1, 2, 3] ← изменился!
# ✅ Хорошо: sorted() создаёт новый список
def good_pipeline(data):
sorted_data = sorted(data)
return [x * 2 for x in sorted_data]
numbers = [3, 1, 2]
result = good_pipeline(numbers)
print(numbers) # [3, 1, 2] ← не изменился
Ошибка 2: Слишком глубокая вложенность
# ❌ Нечитаемо
result = reduce(
lambda a, x: a + x,
map(lambda x: x ** 2,
filter(lambda x: x > 0,
map(lambda x: x - 10,
filter(lambda x: x % 2 == 0, data)))),
0
)
# ✅ Разбить на именованные шаги
step1 = [x for x in data if x % 2 == 0]
step2 = [x - 10 for x in step1]
step3 = [x for x in step2 if x > 0]
result = sum(x ** 2 for x in step3)
Правила хорошего pipeline
- Каждый шаг — чистая функция (не меняет входные данные)
- Сложные pipeline разбивать на именованные шаги
- Используй list comprehension там, где это читаемее
- Для больших данных применяй генераторы вместо списков
💬 Комментарии (0)
Комментариев пока нет
Станьте первым, кто поделится мнением об этой статье!