📝 Python

Data Pipelines — Конвейер обработки данных

P
Автор
PyLand Team
📅
Опубликовано
03.04.2026
⏱️
Время чтения
3 мин
👁️
Просмотров
368
🌳
Уровень
Продвинутый

Конвейер данных последовательно передаёт результат одного этапа следующему и помогает разделить сложную обработку на небольшие операции.

Что такое 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)

🔐 Войдите в систему, чтобы оставить комментарий
🚪 Войти
💭

Комментариев пока нет

Станьте первым, кто поделится мнением об этой статье!

🔗 Похожие

Похожие статьи

Продолжите изучение с этими материалами

📝

if __name__ == "__main__": точка входа Python

Конструкция if name == "main" определяет точку входа Python-программы: она помогает отличить прямой запуск файла...

📅 14.08.2026 👁️ 44
📝

strip() и lower(): подготавливаем пользовательски…

Пользователь может ввести правильное слово с лишними пробелами или буквами другого регистра. Для Python строки...

📅 11.08.2026 👁️ 50
📝

ord(), chr() и циклический сдвиг букв в Python 🔐

Строка состоит из символов, но компьютер хранит каждый символ как числовой код. Python позволяет переходить...

📅 09.08.2026 👁️ 89

Понравилась статья?

Подпишитесь на наши обновления и получайте новые статьи первыми. Развивайтесь вместе с PyLand!