Nexa Intelligence Docs
Pipeline Data

articles

Alur terdistribusi pengisian tabel PostgreSQL articles dan indeks Elasticsearch articles—mulai dari penjadwalan taksonomi A1, perayapan B1/B2, pembersihan C1, deduplikasi vektor C2, hingga ekstraksi AI SPOK D1.

Pipeline Data: articles (Master Berita Siber)

Pipeline ini bertanggung jawab untuk mengisi dan memperbarui dokumen berita siber pada tabel PostgreSQL articles (sebagai Single Source of Truth OLTP) dan indeks Elasticsearch articles (sebagai mesin pencarian teks cepat dan agregasi visualisasi).


1. Ringkasan Arsitektur Alur

[taxonomy_keywords] (PostgreSQL)
        │
        ▼ (Cron Trigger)
[A1 - GolangScraper] ──► Antrian Kafka: queue-job
        │
        ▼
[B1/B2 - Crawler Sites] ──► Antrian Kafka: queue-raw-data-site
        │
        ▼
[C1 - HTML Cleaner] ──► Antrian Kafka: queue-normalize-data
        │
        ▼
[C2 - Vector Deduplication] (Qdrant Cosine Similarity < 0.90)
        │
        ▼ (Artikel Unik) ──► Antrian Kafka: queue-ai-inference
[nexa-d1 - AI Worker] (SPOK, NER, Sentimen, Hierarki Polri via AIOS)
        │
        ├────────────────────────────────────┐
        ▼                                    ▼
[PostgreSQL: articles]             [Elasticsearch: articles]
(Master Transaksional)             (Pencarian Cepat & Dashboard)

2. Tahap 1: Penjadwalan & Distribusi Tugas (A1)

Worker A1 GolangScraper membaca master kata kunci dari PostgreSQL (taxonomy_keywords) dan daftar target portal berita (scraping_targets). Setiap tugas crawling dibungkus sebagai JSON job dan dipancarkan ke antrian Kafka queue-job.

Alur Kerja A1 Scraper Controller

Sinkronisasi taksonomi, penjadwalan perayapan portal, hingga pemancaran antrian Kafka

01. Target & Keyword RegistryPostgreSQL

Penyimpanan master daftar domain media, URL RSS/sitemap, dan taksonomi keyword L1-L5.

02. Cron Job & Frequency SchedulerGolang Cron Engine

Menghitung interval perayapan dinamis per media (Tier 1 tiap 5 menit, Tier 2 tiap 15 menit).

03. Ingestion Task DispatcherTask Generator

Membungkus instruksi perayapan menjadi payload job terstandarisasi dengan ID unik.

04. Kafka Producer Ingestion QueueTopic: queue-job

Publikasi message task ke antrian Kafka untuk dikonsumsi oleh kluster scraper B1.


3. Tahap 2: Pengambilan Konten & Sanitasi Teks (B1, B2 & C1)

  1. B1 & B2 Scraper: Mengambil HTML dokumen dari portal Tier 1 hingga Tier 3, lalu memancarkannya ke Kafka queue-raw-data-site.
  2. C1 Cleaner: Membuang elemen DOM iklan, tracker, header, dan footer boilerplate. Mengonversi waktu publikasi ke format standar ISO-8601 UTC dan memancarkan teks bersih ke Kafka queue-normalize-data.

4. Tahap 3: Deduplikasi Semantik Vektor (C2)

Worker C2 membandingkan representasi vektor judul dan paragraf pembuka artikel ke koleksi Qdrant Vector DB:

  • Skor Kemiripan >= 0.90 (Duplikat / Sindikasi): Ditolak dari pipeline inferensi AI untuk menghemat biaya komputasi, dialihkan ke pencatatan article_trendings.
  • Skor Kemiripan < 0.90 (Artikel Unik / Baru): Dipancarkan ke Kafka queue-ai-inference untuk dianalisis oleh model AI kognitif.

Alur C2 Deduplication Worker

Deduplikasi semantik berbasis embedding vector dan pengelompokan cluster berita

01. Ingest Event & Fetch Article BodyKafka • MinIO Storage

Mengambil metadata artikel dari Kafka dan menarik konten teks lengkap dari MinIO.

02. Generate Dense EmbeddingFastEmbed 384-dim

Vektorisasi judul dan ringkasan artikel menjadi dense vector embedding 384 dimensi.

03. Qdrant Vector Similarity SearchThreshold: ≥ 0.85

Query pencarian ANN pada collection Qdrant untuk mendeteksi kesamaan konten berita.

Duplikat / Sindikasi (≥ 0.85)
Tautkan ke cluster ID existing & increment dupe count.
Berita Unik (< 0.85)
Simpan vector ke Qdrant & teruskan ke antrian AI D1.

5. Tahap 4: Ekstraksi Kognitif & Dual-Sink Persistence (D1)

Worker nexa-d1 mengonsumsi event dari queue-ai-inference:

  1. Ekstraksi Sintaksis SPOK: Membedah kalimat inti berita menjadi Subjek, Predikat, Objek, dan Keterangan.
  2. Polri Hierarchy Resolver: Mengidentifikasi unit satuan kerja kepolisian dan memetakan ke hierarki Mabes $\rightarrow$ Polda $\rightarrow$ Polres $\rightarrow$ Polsek.
  3. Sentimen & Skor Risiko: Menilai polaritas (Positif, Netral, Negatif) dan skor urgensi kamtibmas (1–5).
  4. Penyimpanan Ganda:
    • INSERT INTO articles di PostgreSQL.
    • POST /articles/_doc/{id} di Elasticsearch.

Alur D1 AI Worker (Berita)

Inferensi LLM terdistribusi untuk sentimen, entitas, ringkasan, dan dual-sink

01. Ingestion Task ConsumerKafka: topic-ai-worker

Menerima payload artikel bersih pasca-deduplikasi untuk diproses secara batch asinkron.

02. Multi-Task LLM InferencevLLM • Qwen / Mistral

Ekstraksi simultan: Sentimen (-1 s/d +1), NER (Tokoh, Organisasi, Lokasi), Ringkasan 3-kalimat, & Relevansi Taksonomi.

MongoDB Document Store
Collection: articles_clean_v2
Elasticsearch Search Index
Index: news_search_index
04. Downstream Event DispatchKafka: topic-aggregation • topic-decision

Memicu aggregation worker & decision worker untuk kalkulasi tren dan scoring ancaman otomatis.


6. Spesifikasi Payload Request Pipeline

Berikut adalah contoh payload data yang ditransmisikan pada setiap titik jabat tangan (handshake) antar-worker:

Payload pesan Kafka yang dipancarkan oleh worker C2 setelah artikel diverifikasi unik (skor Cosine < 0.90 di Qdrant):

{
  "job_id": "a98213f4-1234-4567-89ab-cdef01234567",
  "article_id": "c1f7a01e-4501-49b8-a73c-7c093a123456",
  "url": "https://nasional.kompas.com/read/2026/09/25/sidak-beras-bulog",
  "domain": "kompas.com",
  "source_name": "Kompas.com",
  "title": "Satgas Pangan Polri Gelar Sidak Gudang Beras di Jawa Barat",
  "clean_content": "Satgas Pangan Polri bersama Ditreskrimsus Polda Jawa Barat menggelar inspeksi mendadak di sejumlah gudang penyimpanan beras di wilayah Karawang guna memastikan kelancaran distribusi beras SPHP...",
  "published_at": "2026-09-25T08:30:00Z",
  "qdrant_point_id": "c1f7a01e-4501-49b8-a73c-7c093a123456",
  "vector_similarity_score": 0.42
}

On this page