Скользящее Среднее за 3 месяца: Решение с Оконными Функциями
Это классическое задание по работе с временными рядами в SQL. Оконные функции (Window Functions) — идеальный инструмент для этого.
1. Базовое Скользящее Среднее за 3 месяца
SELECT
DATE_TRUNC('month', date)::DATE AS month,
SUM(income) AS monthly_income,
AVG(SUM(income)) OVER (
ORDER BY DATE_TRUNC('month', date)
ROWS BETWEEN 2 PRECEDING AND CURRENT ROW
) AS moving_avg_3m
FROM transactions
GROUP BY DATE_TRUNC('month', date)
ORDER BY month;
Результат:
month | monthly_income | moving_avg_3m
2024-01-01 | 10000.00 | 10000.00
2024-02-01 | 15000.00 | 12500.00
2024-03-01 | 12000.00 | 12333.33
2024-04-01 | 18000.00 | 15000.00
2024-05-01 | 14000.00 | 14666.67
2. С Обработкой Граничных Случаев (NULLs для недостаточных данных)
Решение
Архитектурный подход
Двухслойная архитектура:
Этот подход позволяет:
Слой 1: Hive (Raw Data)
sales_raw — логи всех продажCREATE TABLE IF NOT EXISTS sales_raw (
sale_id BIGINT,
sale_date DATE,
sale_timestamp TIMESTAMP,
store_id INT,
sku_id INT,
category_id INT,
city_name STRING,
quantity INT,
unit_price DECIMAL(10,2),
total_amount DECIMAL(15,2),
cost_amount DECIMAL(15,2),
margin_amount DECIMAL(15,2),
load_date DATE
)
PARTITIONED BY (year INT, month INT)
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\t';
Решение
Анализ задачи
Дубликаты в таблице: Для каждого client_id есть несколько записей с разными updated_at. Нужно оставить только самую свежую запись для каждого клиента.
В примере:
Решение 1: SELECT последних записей (просмотр)
SELECT DISTINCT ON (client_id)
id,
client_id,
balance,
updated_at
FROM ClientBalance
ORDER BY client_id, updated_at DESC;
Только для PostgreSQL. Объяснение:
DISTINCT ON (client_id) — оставляет одну запись на client_idORDER BY client_id, updated_at DESC — первая запись в группе (самая свежая)Результат:
id | client_id | balance | updated_at
---|-----------|---------|---------------------
4 | 100 | 6000.00 | 2024-01-17 12:00:00
3 | 200 | 3000.00 | 2024-01-15 09:00:00
Решение 2: С использованием ROW_NUMBER (универсальное)
Решение
Задача и подход
Нужно определить сессии пользователей с интервалом инактивности менее 30 минут. Для этого используем PySpark с функциями работы с окнами (Window Functions) и обнаружением разрывов в хронологической последовательности.
Основной алгоритм:
Шаг 1: Присвоение session_id
from pyspark.sql import SparkSession, Window
from pyspark.sql.functions import *
from datetime import datetime
spark = SparkSession.builder.appName("SessionAnalysis").getOrCreate()
Решение
Анализ требований
Задача: Найти пользователей, чьи покупки за 3 месяца превышают общее среднее значение.
Ключевые моменты:
Решение 1: С использованием оконной функции AVG (рекомендуется)
Решение
Типовые проблемы масштабирования Spark алгоритмов
Это самая частая проблема. Когда данные распределены неравномерно по партициям, некоторые worker'ы обрабатывают в 100+ раз больше данных, чем другие.
Пример:
# Плохо — ключи распределены неравномерно
df.groupBy("user_id").count().show()
# Если у одного пользователя 99% всех записей — один partition получит все данные
При join'е двух больших таблиц Spark должен переместить данные по сети. Если объём shuffle'а больше доступной памяти, начинаются disk spills (запись на диск), что замедляет выполнение в 10+ раз.
Решение
1. Архитектура системы
Рекомендуемая пятислойная архитектура с разделением ответственности:
SOURCES → INGESTION → PROCESSING → WAREHOUSE → PRESENTATION
Ключевые компоненты:
2. Технологический стек
Ingestion Layer:
Processing Layer:
Storage Layer:
Решение
Задача и контекст
Нужно найти N-ую самую высокую уникальную зарплату. Это важно различать:
В примере уникальные зарплаты: 100000 (Alice, Diana), 90000 (Charlie), 80000 (Bob), 75000 (Eve)
Решение 1: С использованием LIMIT + OFFSET (простое и быстрое)
SELECT DISTINCT salary
FROM Employee
ORDER BY salary DESC
LIMIT 1 OFFSET 1; -- N-1 для получения N-ой зарплаты
Для N = 2:
SELECT DISTINCT salary
FROM Employee
ORDER BY salary DESC
LIMIT 1 OFFSET 1;
-- Результат: 90000
Объяснение:
DISTINCT — исключаем дубликаты зарплатORDER BY salary DESC — сортируем по убываниюLIMIT 1 — берём одну строкуOFFSET N-1 — пропускаем первые N-1 строкПреимущества:
Решение
Задача и контекст
Необходимо реализовать класс для работы с разреженными векторами (sparse vectors) — векторами, содержащими много нулевых элементов. Эта задача критична в обработке данных, так как реальные данные часто имеют высокую разреженность (тексты, пользовательские рейтинги, сетевые графы). Наивное хранение всех элементов может привести к огромным потерям памяти.
Базовое решение
Начнём с простого подхода, используя встроенные структуры Python:
class SparseVector:
def __init__(self, nums: list[int]):
self.nums = nums
def dotProduct(self, vec: "SparseVector") -> int:
return sum(a * b for a, b in zip(self.nums, vec.nums))
# Пример использования
nums1 = [1, 2, 0, 4]
nums2 = [8, 0, 3, 5]
vec1 = SparseVector(nums1)
vec2 = SparseVector(nums2)
print(vec1.dotProduct(vec2)) # 28
Решение
Задача и подход
Нужно обнаружить непрерывные периоды активности (консецутивные дни) для каждого пользователя. Используем window functions для вычисления разницы между текущей датой и рангом, чтобы идентифицировать разрывы.
Основной алгоритм:
Шаг 1: Определение непрерывных периодов
Решение
Задача и подход
Нужно классифицировать узлы дерева по трём типам: Root (корень), Leaf (лист), Inner (внутренний). Используем стандартные SQL техники с самоприсоединением (self-join) для определения наличия дочерних узлов и анализа родительских связей.
Основной алгоритм:
Решение на SQL
SELECT
t.node_id,
CASE
WHEN t.parent_id IS NULL THEN 'Root'
WHEN NOT EXISTS (
SELECT 1
FROM tree children
WHERE children.parent_id = t.node_id
) THEN 'Leaf'
ELSE 'Inner'
END AS node_type
FROM tree t
ORDER BY t.node_id;
Альтернативный подход с LEFT JOIN
Решение
Понимание медианы
Медиана — значение, которое делит отсортированный набор пополам:
В примере:
Решение 1: С использованием PERCENTILE_CONT (рекомендуется)
Это встроенная функция в большинстве БД для расчёта медианы:
SELECT
department,
PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY salary) AS median_salary
FROM employees
GROUP BY department
ORDER BY department;
Объяснение:
PERCENTILE_CONT(0.5) — 50-й процентиль (медиана)WITHIN GROUP (ORDER BY salary) — сортировка значений для расчётаGROUP BY department — вычисляем медиану в каждой группеРешение
SQL-запрос для подсчёта уникальных значений
SELECT COUNT(DISTINCT user_id) AS unique_users
FROM large_dataset;
Альтернативная форма:
SELECT COUNT(*)
FROM (
SELECT DISTINCT user_id
FROM large_dataset
) t;
Детальное объяснение выполнения в Spark
SparkSQL создаёт логический план:
Aggregate [COUNT(DISTINCT user_id)]
└── Scan parquet large_dataset
HashAggregate (final)
├── Exchange hashpartitioning(user_id, 200) ← Shuffle происходит здесь
└── HashAggregate (partial)
└── Scan parquet large_dataset
Что происходит:
Решение
Проблемы в коде с партиционированием по double
Этот код содержит критическую ошибку, которая может полностью вывести кластер из строя:
Партиционирование по полю amount типа double создаёт отдельную папку для каждого уникального значения. С диапазоном от 0 до 1 млрд значений получится миллиарды папок в HDFS, каждая под отдельное значение.
Если в датасете миллиард уникальных значений — получим миллиард папок!
NameNode хранит в памяти весь namespace HDFS — информацию о каждом файле, папке, блоке данных. На каждую папку требуется примерно 150-250 байт памяти. Для 1 млрд папок это:
1,000,000,000 × 200 байт = 200 ГБ памяти
Типичный NameNode имеет 16-64 ГБ памяти — этого недостаточно. NameNode начнёт работать с диском, что приведёт к полной неработоспособности кластера, GC pauses на десятки секунд, невозможности выполнять какие-либо операции.
Роль NameNode
Решение
Задача и подход
Нужно анализировать переводы между аккаунтами: рассчитать изменение баланса, определить положительные/отрицательные балансы и найти наиболее активные пары аккаунтов.
Шаг 1: Изменение баланса для каждого аккаунта
WITH account_cash_flow AS (
SELECT
account_id,
SUM(CASE WHEN flow_type = 'sent' THEN amount ELSE 0 END) AS total_sent,
SUM(CASE WHEN flow_type = 'received' THEN amount ELSE 0 END) AS total_received
FROM (
SELECT sender_id AS account_id, amount, 'sent' AS flow_type
FROM transactions
UNION ALL
SELECT receiver_id, amount, 'received'
FROM transactions
) flows
GROUP BY account_id
)
SELECT
account_id,
total_sent,
total_received,
total_sent - total_received AS balance_change
FROM account_cash_flow
ORDER BY balance_change DESC;
Результат: 100: -350, 200: +200, 300: +150
Шаг 2: Классификация на положительные и отрицательные
Решение
Анализ задачи
Нужно найти пары (X₁, Y₁) и (X₂, Y₂) из одной таблицы, где:
В примере:
Решение 1: Базовое (использование JOIN)
SELECT DISTINCT
p1.X,
p1.Y
FROM pairs p1
JOIN pairs p2 ON p1.X = p2.Y AND p1.Y = p2.X
WHERE p1.X <= p1.Y -- Исключаем дубликаты (выводим только меньшее значение первым)
ORDER BY p1.X ASC;
Объяснение:
JOIN: Соединяем таблицу с самой собой (p1 JOIN p2)
p1.X = p2.Y AND p1.Y = p2.X — условие симметричностиWHERE p1.X <= p1.Y: Фильтруем, чтобы вывести только одну из двух симметричных пар