Files
2018-lpon-site/_prj_blueprint/BACKGROUND-TASKS-CONCEPT.md
T

21 KiB
Raw Blame History

ПОТОК ДАННЫХ (чистая архитектура) ДРАФТ

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 очередь и кликает:

  1. " Одобрить"job:create:style:1 → Воркер пишет в БД
  2. "🔗 Мержить с ID:5"job:merge:style:1:with:5 → Воркер обновляет existing
  3. " Пропустить" → Удалить из Redis (ничего не пишется)
  4. "📝 Отредактировать" → Изменить 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           │
        └─────────────────────┘