Kafka
Panduan komprehensif seluruh topik Apache Kafka, wire contract, skema JSON, consumer groups, mekanisme DLQ, dan monitoring lag pada platform Nexa Intelligence.
Spesifikasi Topik Kafka & Event Message Bus
Apache Kafka bertindak sebagai backbone asynchronous event-driven streaming di dalam ekosistem Nexa Intelligence Platform. Seluruh komunikasi antar-layanan (mulai dari A1 Scheduler, B1 Scrapers, C1 Normalizer, C2 Deduplication, D1 AIOS, hingga F1 API Backend) terikat oleh kontrak data ketat yang didefinisikan dalam repository contracts-registry.
Interactive Topology & Monitoring Lag Flow
Diagram interaktif berikut memvisualisasikan alur pipeline pesan Kafka, laju throughput messages/sec, titik potensi consumer lag, dan mekanisme Dead-Letter Queue (DLQ):
Alur Pemantauan Consumer Lag Kafka
Pengukuran throughput, offset tracking, deteksi bottleneck antrian, dan alerting
Produser memancarkan pesan secara berkelanjutan ke berbagai partisi topik Kafka.
Menyimpan buffer pesan terdistribusi dengan penanda posisi indeks offset pesan terbaru.
Menghitung selisih (lag) antara pesan yang masuk dan kecepatan konsumsi worker D1 AI dan consumer downstream.
Memicu alert otomatis ke DevOps jika terjadi penumpukan antrian akibat penurunan performa worker.
Matriks Lengkap Topik Kafka Nexa
Berikut adalah matriks seluruh topik Kafka aktif, produsen, konsumen, consumer group, dan strategi partisi:
| Topik Kafka | Bus Contract | Producer (Penerbit) | Consumer (Konsumen) | Consumer Group ID | Partisi & Key Strategy |
|---|---|---|---|---|---|
queue-job | jobbus | nexa-a (A1 Scheduler) | nexa-b1 (Scraper Worker) | b1-scraper | Key: job_id (UUID) |
queue-raw-data-site | rawbus | nexa-b1 (Scraper Worker) | nexa-c1 (Normalizer/Cleaner) | c1-processor | Key: externalId / URL hash |
queue-normalize-data | articlebus | nexa-c1 (Normalizer) | nexa-c2 (Dedup Worker) | c2-dedup-group | Key: article_id (UUID) |
queue-normalize-data-dlq | dlq | nexa-c1, nexa-c2, nexa-d1 | f1-api-backend (Dead-Jobs) | f1-deadjobs-monitor | Key: origin_key |
queue-ai-inference | inference | nexa-c2 (Dedup Worker) | nexa-d1 (AIOS Worker) | d1-inference-group | Key: job_id (UUID) |
queue-trending-data | trending | nexa-c2 (Dedup Worker) | nexa-d1 (Trending Thread) | d1-trending-group | Key: job_id (UUID) |
Redis Streams Legacy Support
Topik scrape-tasks (taskbus) menggunakan transport Redis Streams (XADD / XREADGROUP / XAUTOCLAIM) untuk mendukung scraper legacy sebelum migrasi penuh ke Kafka queue-job.
1. Topik queue-job (A1 → B1)
Topik ini mentransmisikan instruksi crawling dan scraping terjadwal dari A1 Scraper Controller ke pool worker B1 Scrapers.
Spesifikasi Wire
- Transport: Kafka
- Format: JSON (UTF-8)
- Key:
job_id(UUIDv4) - Schema ID:
https://nexa.app/contracts/kafka/queue-job.schema.json
{
"job_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
"triggered_at": "2026-09-10T14:30:00Z",
"target_id": "a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11",
"domain_name": "kompas.com",
"base_url": "https://nasional.kompas.com",
"sector_id": "b3f07a7e-1a1a-4d66-88ef-5883d6a4c211",
"sector_name": "Politik, Hukum, dan Keamanan",
"region": {
"region_id": "32000000-0000-0000-0000-000000000000",
"region_name": "Jawa Barat",
"region_level": "PROVINSI"
},
"organization": {
"org_id": "c1f10a8e-2b2b-4e77-99ff-6994e7b5d322",
"org_name": "Dinas Komunikasi dan Informatika",
"org_type": "PEMDA",
"parent_org_id": null
},
"scraping_config": {
"engine": "chromium_stealth",
"max_pages": 15,
"timeout_ms": 30000,
"html_selectors": {
"article_wrapper": "div.read__content",
"title": "h1.read__title",
"body_content": "div.clearfix p",
"published_date": "div.read__time"
}
},
"lexicon_matcher": [
{
"keyword_id": "d2f21b9e-3c3c-4f88-00aa-7aa5f8c6e433",
"phrase": "pilkada serentak",
"phrases": ["pilkada serentak", "pemilihan kepala daerah", "kpu jabar"],
"taxonomy_lineage": {
"leaf_node_id": "e3f32c0f-4d4d-4a99-11bb-8bb6a9d7f544",
"leaf_node_name": "Pilkada & Elektoral",
"node_level": 3,
"ancestry_path_ids": [
"b3f07a7e-1a1a-4d66-88ef-5883d6a4c211",
"f4a43d10-5e5e-4baa-22cc-9cc7b0e8a655",
"e3f32c0f-4d4d-4a99-11bb-8bb6a9d7f544"
],
"ancestry_path_names": [
"Politik, Hukum, dan Keamanan",
"Demokrasi & Tata Kelola",
"Pilkada & Elektoral"
]
}
}
]
}2. Topik queue-raw-data-site (B1 → C1)
Topik ini memuat hasil scraping mentah yang dihasilkan oleh worker B1 Scraper sebelum dibersihkan (sanitized) oleh C1 Normalizer.
Spesifikasi Wire
- Transport: Kafka
- Key:
externalId(URL hash MD5/SHA256) - Schema ID:
https://nexa.app/contracts/kafka/queue-raw-data-site.schema.json
{
"job_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
"target_id": "a0eebc99-9c0b-4ef8-bb6d-6bb9bd380a11",
"externalId": "e4d909c290d0fb1ca068ffaddf22cbd0",
"url": "https://nasional.kompas.com/read/2026/09/10/kpk-periksa-saksi-lelang",
"domain_name": "kompas.com",
"raw_html": "<!DOCTYPE html><html><head>...",
"minio_s3_key": "raw/news/kompas.com/2026-09-10/e4d909c290d0fb1ca068ffaddf22cbd0.json",
"scraped_at": "2026-09-10T14:35:12Z",
"http_status": 200,
"crawler_metadata": {
"engine": "chromium_stealth",
"ip_proxy": "103.145.22.10",
"duration_ms": 1420
}
}3. Topik queue-normalize-data (C1 → C2)
Topik ini memuat data artikel yang telah dibersihkan dari tag HTML, dinormalisasi stempel waktunya, dan siap diuji deduplikasi semantik oleh C2 Deduplication Engine.
Spesifikasi Wire
- Transport: Kafka
- Key:
article_id(UUIDv4) - Schema ID:
https://nexa.app/contracts/kafka/queue-normalize-data.schema.json
{
"job_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
"article_id": "a1b2c3d4-e5f6-7a8b-9c0d-1e2f3a4b5c6d",
"url": "https://nasional.kompas.com/read/2026/09/10/kpk-periksa-saksi-lelang",
"domain": "kompas.com",
"title": "KPK Periksa Saksi Terkait Dugaan Korupsi Proyek Smart City",
"author": "Fauzi Rahman",
"content": "Jakarta - Komisi Pemberantasan Korupsi (KPK) menjadwalkan pemeriksaan terhadap tiga orang saksi...",
"published_at": "2026-09-10T14:00:00Z",
"cleaned_at": "2026-09-10T14:36:01Z",
"char_count": 2850,
"word_count": 420,
"language": "id",
"matched_keywords": ["kpk", "korupsi", "smart city", "pengadaan barang"],
"region": {
"region_id": "32000000-0000-0000-0000-000000000000",
"region_name": "Jawa Barat"
},
"sector": {
"sector_id": "b3f07a7e-1a1a-4d66-88ef-5883d6a4c211",
"sector_name": "Politik, Hukum, dan Keamanan"
}
}4. Topik queue-ai-inference (C2 → D1)
Topik ini mentransmisikan artikel yang telah lolos deduplikasi semantik (artikel unik) ke nexa-d1 untuk diekstraksi sintaksis SPOK, relasi entitas, normalisasi hierarki lembaga, dan scoring risiko pengadaan.
Spesifikasi Wire
- Transport: Kafka
- Key:
job_id(UUIDv4) - Schema ID:
https://nexa.app/contracts/kafka/queue-ai-inference.schema.json
{
"job_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
"article_id": "a1b2c3d4-e5f6-7a8b-9c0d-1e2f3a4b5c6d",
"title": "KPK Periksa Saksi Terkait Dugaan Korupsi Proyek Smart City",
"content": "Jakarta - Komisi Pemberantasan Korupsi (KPK) menjadwalkan pemeriksaan terhadap tiga orang saksi...",
"domain": "kompas.com",
"published_at": "2026-09-10T14:00:00Z",
"qdrant_vector_id": "c1f7a01e-4501-49b8-a73c-7c093a123456",
"dedup_status": "UNIQUE",
"similarity_score": 0.32,
"task_priority": "HIGH"
}5. Topik queue-trending-data (C2 → D1 Trending)
Topik ini memuat artikel yang terdeteksi sebagai Duplikat / Sindikasi (Cosine Similarity ≥ 0.90). Artikel ini tidak melewati inferensi AI yang mahal, melainkan langsung dialirkan untuk menambah bobot kecepatan (velocity) dan tren sebutan isu.
Spesifikasi Wire
- Transport: Kafka
- Key:
original_article_id(UUID artikel master) - Schema ID:
https://nexa.app/contracts/kafka/queue-trending-data.schema.json
{
"duplicate_article_id": "f5e4d3c2-b1a0-9876-5432-10fedcba9876",
"original_article_id": "a1b2c3d4-e5f6-7a8b-9c0d-1e2f3a4b5c6d",
"domain": "tribunnews.com",
"published_at": "2026-09-10T14:45:00Z",
"similarity_score": 0.94,
"duplicate_cluster_id": "cluster_smart_city_kpk_20260910"
}6. Topik queue-normalize-data-dlq (Dead-Letter Queue)
Topik ini menampung seluruh pesan yang mengalami kegagalan ekstraksi atau pemrosesan setelah melewati batas maksimal percobaan (retry exhaustion).
Pesan di topik DLQ dipantau secara real-time oleh f1-api-backend melalui modul Dead-Jobs Monitor untuk investigasi kegagalan dan dapat di-replay langsung via API atau Dashboard.
Format Dokumen DLQ:
{
"dlq_id": "dlq_8921739812",
"origin_topic": "queue-raw-data-site",
"failed_stage": "nexa-c1-normalizer",
"error_code": "HTML_PARSING_FAILED",
"error_message": "Selector div.read__content not found in DOM structure",
"attempt_count": 3,
"first_failed_at": "2026-09-10T14:35:15Z",
"last_failed_at": "2026-09-10T14:35:45Z",
"is_replayable": true,
"original_key": "e4d909c290d0fb1ca068ffaddf22cbd0",
"original_payload": {
"job_id": "9b1deb4d-3b7d-4bad-9bdd-2b0d7b3dcb6d",
"url": "https://nasional.kompas.com/read/2026/09/10/kpk-periksa-saksi-lelang"
}
}knowledge_graph
Alur pembentukan jejaring graf pengetahuan entitas aktor-organisasi-isu harian menggunakan AIOS Dify Workflow, kalkulasi sentralitas graf, dan penyimpanan ke indeks Elasticsearch knowledge_graph serta PostgreSQL knowledge_graphs.
PostgreSQL
Spesifikasi komprehensif skema database relasional OLTP PostgreSQL, definisi tabel master, relasi entitas, struktur artikel AI, data sektoral, serta strategi indeks performa di Nexa Intelligence Platform.