Знакомая ситуация: нужно выгрузить данные за 30 дней, пишешь цикл, день за днём, и ждёшь 30 минут.
А можно запустить все запросы одновременно и получить результат в 5–10 раз быстрее.
Сегодня покажу, как это сделать на простом примере.
Что сделаем:
🟢 Напишем функцию, которая выгружает данные за один день из ClickHouse.
🟢 Запустим её для нескольких дат параллельно с помощью
ThreadPoolExecutor.🟢Соберём все результаты в один
DataFrame.Разбор по шагам:
1️⃣ Функция для одного дня:
Создаём функцию
load_day — она принимает дату, делает запрос в ClickHouse и возвращает DataFrame:import clickhouse_connect
import pandas as pd
def load_day(date_str):
print(f' [{date_str}] → запрос начат')
client = clickhouse_connect.get_client(**CH_CREDS)
df = client.query_df(f"""
SELECT * FROM events
WHERE date = '{date_str}'
""")
client.close()
print(f' [{date_str}] → готово, строк: {len(df)}')
return df
2️⃣ Запускаем параллельно:
Используем
ThreadPoolExecutor. Ключевой момент — функция начинает выполняться прямо в момент вызова submit():from concurrent.futures import ThreadPoolExecutor, as_completed
dates = ['2024-01-01', '2024-01-02', '2024-01-03', '2024-01-04']
executor = ThreadPoolExecutor(max_workers=3)
tasks = {executor.submit(load_day, d): d for d in dates}
# ↑ вот здесь load_day УЖЕ запустилась в фоновом потоке
# submit() не ждёт результата — сразу возвращает "квитанцию" (Future)
# и цикл переходит к следующей дате
Что происходит под капотом:
•
submit(load_day, '2024-01-01') → поток 1 сразу побежал делать запрос в ClickHouse •
submit(load_day, '2024-01-02') → поток 2 сразу побежал, пока первый ещё работает •
submit(load_day, '2024-01-03') → поток 3 сразу побежал •
submit(load_day, '2024-01-04') → все 3 воркера заняты — ждёт в очередиТо есть
submit() — это как «отдать задание курьеру». Курьер побежал, а ты уже отдаёшь задание следующему. max_workers=3 — значит у тебя 3 курьера. Четвёртое задание ляжет на стол и подождёт, пока кто-то вернётся.3️⃣ Собираем результаты:
as_completed отдаёт задачи по мере их завершения — какая первая закончилась, ту и получаем:results = []
for task in as_completed(tasks):
date = tasks[task]
try:
results.append(task.result())
print(f' {date} — результат получен ✅')
except Exception as e:
print(f' {date} — ошибка: {e} ❌')
executor.shutdown()
df = pd.concat(results, ignore_index=True)
print(f'Всего строк: {len(df)}')
🌟 Как это выглядит в консоли:
[2024-01-01] → запрос начат ← воркер 1 стартовал
[2024-01-02] → запрос начат ← воркер 2 стартовал
[2024-01-03] → запрос начат ← воркер 3 стартовал
[2024-01-03] → готово, строк: 1200
2024-01-03 — результат получен ✅
[2024-01-04] → запрос начат ← воркер 3 освободился, взял 4-й день
[2024-01-01] → готово, строк: 980
2024-01-01 — результат получен ✅
...
Итог: 🤩
Как выбирать max_workers?
• Для тяжёлых запросов к БД: 3–5 воркеров
• Для лёгких API-запросов: 10–20 воркеров
• Начинайте с 3 и смотрите на нагрузку базы — ClickHouse может обрабатывать ограниченное число запросов одновременно
🚀 Работает не только с ClickHouse — PostgreSQL, MySQL, любые API. Паттерн универсальный.
🔥 Набираем 70 реакций на этот пост — так я пойму, что стоит делать больше таких разборов.
📤 Пересылай пост коллеге или другу из аналитики — это бесплатно и помогает каналу расти. Спасибо!
🚬 Вопросы, менторство, консультации: Написать в ЛС | mentor.dima-sqlit.ru
@dima_sqlit
