21 KiB
21 KiB
ПОТОК ДАННЫХ (чистая архитектура) ДРАФТ
1️⃣ ИСТОЧНИК ДОБАВЛЕН
├─ Пользователь добавляет новый источник (TbSource)
└─ Создается задание в Redis: `job:parse_source:123`
2️⃣ ФОНОВЫЙ ВОРКЕР (Celery/APScheduler)
├─ Видит задание в Redis
├─ Начинает парсить (не пишет в prod БД!)
├─ Получает: Style["Rock"], Artist["The Beatles"], Label["Sony"]
├─ Кидает КАЖДЫЙ в Redis как **виртуальный объект**
│ └─ `pending:style:1` → {"title": "Rock", "source": "Discogs", ...}
│ └─ `pending:artist:2` → {"title": "The Beatles", "source": "MusicBrainz", ...}
└─ Смотрит: есть ли в product БД? Если НЕТ → в Redis очередь
3️⃣ АДМИНКА ДЖАНГО (виртуальная вкладка "Ожидающие задачи")
├─ Админ видит в Redis список
├─ Видит: Rock (похож на existing Rock?), The Beatles (новый?), Sony (уже есть)
├─ **ПОТЫКАЛ**:
│ ├─ Rock → "Мержить с existing Rock" → задание в Redis
│ ├─ The Beatles → "Одобрить" → задание в Redis
│ └─ Sony → "Пропустить" → удалить из Redis
└─ Каждое действие → новое задание в Redis
4️⃣ ФОНОВЫЙ ВОРКЕР II (финальный штрих)
├─ Видит задание: "merge Rock с ID:5"
│ └─ Делает: добавляет синонимы, сохраняет в product БД
├─ Видит задание: "create The Beatles"
│ └─ Пишет вреальную БД → TbArtist создана
└─ После каждого успеха → **удалить из Redis**
5️⃣ ПРОДАКТ БД (в итоге)
└─ Только валидные, одобренные, обработанные данные
└─ No garbage, ID подряд
Архитектура
Парсер Redis (очередь) Админка Воркер БД
│ │ │ │ │
├──parse_source───────────→│ │ │ │
│ │ │ │ │
│ ┌─pending:style:1 │ │ │
│ ├─pending:artist:1 ←─ Админ видит видит здесь │ │
│ └─pending:label:1 потыкает кнопки │ │
│ │ │ │ │
│ │ "merge Rock"│ │ │
│ │←───────────────────│ →─ job:merge:1 ──────→│ │
│ │ │ ├─→ UPDATE Style →│
│ │ │ │ delete from Redis
│ │ "create Beatles" │ │
│ │←──────────────────→─ job:create:2 ────────→│ │
│ │ │ ├─→ INSERT Artist ─│
│ │ │ │ delete from Redis
│ │ │ │ │
└───────────────────────────────────────────────────────────────────────────────────────────────────→
КЛЮЧЕВЫЕ ПРИНЦИПЫ
- Чистая БД: product база получает только одобренные, валидные данные
- ID подряд: мнимизирукет удаления (delete), INSERT данных парсинга только при одобрении
- Single Source of Truth: пока не одобрено администратором → данные ТОЛЬКО в Redis
- Асинхронность: парсер не блокирует админку, админка не блокирует парсер
- Откат дешевый: удалить из Redis дешевле, чем восстанавливать из БД
КОМПОНЕНТЫ REDIS ОЧЕРЕДИ
tasks:pending ← Очередь неодобренных задач (парсер → сюда)
└─ pending:style:1
└─ pending:artist:2
└─ pending:label:3
tasks:approved ← Очередь одобренных (админ → сюда)
└─ {type: 'create', id: 'style:1', data: {...}}
└─ {type: 'merge', id: 'artist:1', merge_with_id: 5}
└─ {type: 'skip', id: 'label:3'}
tasks:completed ← История завершённых (воркер → сюда)
tasks:failed ← История ошибок (воркер → сюда)
ДЕЙСТВИЯ В АДМИНКЕ (виртуальная вкладка)
Админ видит Redis очередь и кликает:
- "✅ Одобрить" →
job:create:style:1→ Воркер пишет в БД - "🔗 Мержить с ID:5" →
job:merge:style:1:with:5→ Воркер обновляет existing - "❌ Пропустить" → Удалить из Redis (ничего не пишется)
- "📝 Отредактировать" → Изменить JSON в оптимистичной форме → сохранить как новое задание
ПАРСЕР (упрощённо)
def parse_source(source_id):
source = TbSource.objects.get(id=source_id)
for style_name in PARSED_STYLES:
# 1. Ищем существующий стиль или по названию или по алиасам
existing = TbMusicStyle.objects.filter(
Q(s_style_name__iexact=style_name) |
Q(j_style_synonyms__contains=style_name)
).first()
if existing:
# ✅ МАТЧИНГ СРАБОТАЛ → пишем сразу в БД
if style_name not in existing.j_style_synonyms:
existing.j_style_synonyms.append(style_name) # Запомнили синоним!
existing.save()
else:
# ❌ НОВЫЙ СТИЛЬ → в Redis очередь на одобрение
redis.lpush('tasks:pending', {...})
# В метаданных источника сохраняем прогресс парсинга
source.j_source_metadata['last_parsed_line'] = current_row_number
source.j_source_metadata['total_lines'] = total_rows
source.j_source_metadata['parsed_at'] = timezone.now().isoformat()
source.save()
АДМИНКА (упрощённо)
class PendingTaskAdmin(admin.ModelAdmin):
# Не наследуем от ModelAdmin (нет моделей!)
# Вместо этого: читаем Redis, рендерим как таблицу
def changelist_view(self, request):
# Получаем из Redis
tasks = [json.loads(redis.get(key)) for key in redis.keys('pending:*')]
# Выводим нетипичную таблицу с кнопками: Одобрить, Мержить, Пропустить
ФОНОВЫЙ ВОРКЕР (упрощённо)
def process_approved_tasks():
while True:
task = redis.lpop('tasks:approved')
try:
if task['type'] == 'create':
# Пишем в product БД
TbMusicStyle.objects.create(title=task['title'], ...)
elif task['type'] == 'merge':
# Обновляем existing + добавляем синонимы
style = TbMusicStyle.objects.get(id=task['merge_with_id'])
style.j_style_synonyms.append(task['title'])
style.save()
# Успех → удалить из Redis
redis.delete(task['_redis_key'])
redis.lpush('tasks:completed', task)
except Exception as e:
redis.lpush('tasks:failed', {**task, 'error': str(e)})
ТЕХНОЛОГИЧЕСКИЙ СТЕК
- Redis — очередь + кэш виртуальных объектов
- Celery или APScheduler — фоновый воркер для парсера и финального сохранения
- Django Admin — кастомная вкладка для управления очередью
- PostgreSQL/SQLite — product database (только одобренные данные)
ВЕТКА: ПРЯМОЕ ПОПАДАНИЕ В БД (когда матчинг сработал)
Если парсер нашел алиас или группа уже в БД → обходим очередь, пишем сразу:
def parse_source(source_id):
source = TbSource.objects.get(id=source_id)
for style_name in PARSED_STYLES:
# 1. Ищем существующий стиль или по названию или по алиасам
existing = TbMusicStyle.objects.filter(
Q(s_style_name__iexact=style_name) |
Q(j_style_synonyms__contains=style_name)
).first()
if existing:
# ✅ МАТЧИНГ СРАБОТАЛ → пишем сразу в БД
if style_name not in existing.j_style_synonyms:
existing.j_style_synonyms.append(style_name) # Запомнили синоним!
existing.save()
else:
# ❌ НОВЫЙ СТИЛЬ → в Redis очередь на одобрение
redis.lpush('tasks:pending', {...})
# В метаданных источника сохраняем прогресс парсинга
source.j_source_metadata['last_parsed_line'] = current_row_number
source.j_source_metadata['total_lines'] = total_rows
source.j_source_metadata['parsed_at'] = timezone.now().isoformat()
source.save()
ИЗМЕНЕНИЕ ЦЕН (дополнительная ветка)
Когда уже существующий товар поменял цену → в отдельную очередь:
tasks:price_changes ← Отдельная очередь
└─ {type: 'price_update', offer_id: 42, old_price: 100, new_price: 120}
└─ {type: 'price_update', offer_id: 43, old_price: 200, new_price: 180}
Обработка:
- Фоновый воркер видит изменение цены
- Проверяет: изменилась ли существенно (> 5%)?
- Если да → может требоваться одобрение (+1 задание в Redis)
- Если нет → пишет сразу в БД (TbOfferHistory записывается автоматически)
МЕТАДАННЫЕ ОТСЛЕЖИВАНИЯ ПРОГРЕССА
В каждой TbSource должно быть поле j_source_metadata:
j_source_metadata = JSONField(default=dict, help_text='Отслеживание парсинга и синонимы')
# Структура:
{
"last_parsed_line": 4523, # Докуда добежал парсер (для resume)
"total_lines": 10000, # Всего строк в файле
"parsed_at": "2026-06-14T12:30:00",# Время последнего парсинга
"status": "in_progress", # in_progress | completed | failed
"error_message": null, # если failed, чтобы видно было почему
# Найденные синонимы (стили, артисты, которые auto-matched)
"matched_styles": {
"Rock": 45, # Стиль 45 найден под названием "Rock"
"rock": 45, # Вариант написания também сохранили
},
"matched_artists": {
"The Beatles": 12,
"Beatles, The": 12,
}
}
Если парсер упал на строке 4523 → перезапуск продолжится с 4524, а не с начала!
МАСШТАБИРУЕМОСТЬ: REDIS В ПАМЯТИ vs PERSISTENCE
Проблема: 10k+ задач при загрузке большого Excel
Excel с 10,000 позиций
├─ 10k артистов
├─ 5k стилей
├─ 2k лейблов
└─ 15k общих задач в очереди
Redis хранит в памяти по умолчанию ➜ контейнер перестартует ➜ всё теряется!
Решение: RDB + AOF Persistence
# docker-compose.yml для Redis
redis:
image: redis:7-alpine
volumes:
- redis-data:/data
command: >
redis-server
--appendonly yes
--appendfsync everysec
--save 900 1
--maxmemory 2gb
--maxmemory-policy allkeys-lru
- RDB snapshots — снимок каждые 15 минут
- AOF log — каждая команда лог записывается на диск
- maxmemory-policy — если память 2GB переполнится, удаляются старые задачи (но только pending, не approved!)
- Persistence: при рестарте контейнера Redis восстановит все задачи из AOF
ПАРАЛЛЕЛЬНАЯ ОБРАБОТКА (масштабирование)
Если задач много → несколько воркеров, каждый обрабатывает свой тип:
# Worker 1: Парсит (пишет в tasks:pending)
celery_app.send_task('parser.parse_source', args=[source_id])
# Worker 2: Обрабатывает стили (слушает tasks:approved тип='style')
@app.task
def process_style_task(task_data):
# обновить БД
# Worker 3: Обрабатывает артистов (слушает tasks:approved тип='artist')
@app.task
def process_artist_task(task_data):
# обновить БД
# Worker 4: Обрабатывает цены (слушает tasks:price_changes)
@app.task
def update_offer_price(task_data):
# обновить цену
Каждый воркер работает в отдельном потоке/процессе ➜ параллелизм!
FLOW ПРИНЦИПИАЛЬНАЯ СХЕМА (обновлённая)
┌──────────────────────────────────┐
│ ПАРСЕР (Worker 1) │
│ Читает Excel/CSV/JSON │
└──────────────┬───────────────────┘
│
┌──────────────┴────────────┐
│ │
┌───────▼───────────┐ ┌─────────▼──────────┐
│ Матч сработал? │ │ Новые данные? │
│ (aliasing) │ │ (неизвестны) │
└────┬──────────────┘ └─────────┬──────────┘
│ Yes │ No
│ │
┌────────▼────────────┐ ┌────────▼────────────┐
│ СРАЗУ В БД │ │ tasks:pending │
│ + j_style_synonyms │ │ (очередь нужнофала) │
│ + metadata['matched'] │ (ждут одобрения) │
└────────┬────────────┘ └────────┬────────────┘
│ │
│ АДМИНКА ДЖАНГО
│ (виртуальная таблица)
│ │
│ ┌──────┴──────┐
│ │ ✅ ❌ 🔗 │
│ (Approve/Skip/Merge)
│ │ │
│ ┌───────▼─────┐ │
│ │ tasks: │ │
│ │ approved │ │
│ └───────┬─────┘ │
│ │ │
└──────────┬─────────┴─────────────┘
│
Worker 2,3,4...
│
┌──────▼──────┐
│ PRODUCT БД │
│ (валидные) │
└─────────────┘
ДОПОЛНИТЕЛЬНЫЕ МЕТАДАННЫЕ ДЛЯ КАЖДОЙ ЗАДАЧИ
# Каждая задача в Redis содержит:
task = {
'id': 'pending:style:123',
'type': 'style', # или 'artist', 'label', 'price_change'
'data': {'title': 'Rock', ...},
# Откуда пришла
'source_id': 42,
'source_name': 'Discogs API',
'parsed_line': 4523,
# Когда создана
'created_at': '2026-06-14T12:00:00',
'ttl_seconds': 86400, # жить 24 часа, потом удалить
# Статус обработки
'attempts': 0, # сколько раз пытались обработать
'last_error': null,
}
ИТОГОВАЯ АРХИТЕКТУРА (расширенная)
ИСТОЧНИКИ (Excel/API)
↓ (15k записей)
┌──────────────────────────────────────────────────┐
│ REDIS ОЧЕРЕДИ (в памяти + AOF persistence) │
│ ├─ tasks:pending (очередь неодобренных) │
│ ├─ tasks:approved (одобренные) │
│ ├─ tasks:price_changes (изменения цен) │
│ ├─ tasks:matched (уже в БД + синонимы) │
│ ├─ tasks:completed (завершённые) │
│ └─ tasks:failed (ошибки) │
└──────────────────────────────────────────────────┘
↓ (видит админка) ↓ (видят воркеры)
┌──────────────────┐ ┌────────────────────┐
│ АДМИНКА ДЖАНГО │ │ WORKERS (Celery) │
│ (виртуальная │ │ ├─ parse_source │
│ таблица) │ │ ├─ process_style │
│ Действия: │ │ ├─ process_artist │
│ ok/del/edit/... │ │ └─ update_prices │
└──────────────────┘ └────────────────────┘
↓ ↓
└──────────┬─────────────┘
↓
┌─────────────────────┐
│ PRODUCT БД │
│ (только валидные) │
│ + metadata progress │
│ + aliases │
└─────────────────────┘