Files
flowfish/docs/analysis/INTEGRATION_STATUS.md
taylanbakircioglu d7ca50b387 feat: Flowfish - eBPF-powered multi-cluster Kubernetes observability platform
Multi-cluster dependency mapping, real-time network monitoring,
impact analysis, and CI/CD integration capabilities.

Made-with: Cursor
2026-03-29 21:43:20 +03:00

5.8 KiB
Raw Permalink Blame History

🐟 Flowfish - Integration Status

RabbitMQ Entegrasyonu

Tarih: 21 Kasım 2024
Durum: Başarıyla entegre edildi

Yapılan İşler

1. Mimari Tasarım

  • Data flow tasarımı (Inspektor Gadget → Ingestion → RabbitMQ → Writer → ClickHouse)
  • Message queue architecture
  • Exchange ve queue konfigürasyonları
  • Message formats (JSON)

2. Kubernetes Deployment

  • 06-rabbitmq.yaml oluşturuldu
  • StatefulSet (persistent storage)
  • ConfigMap (RabbitMQ configuration)
  • Secret (credentials)
  • Service (ClusterIP)
  • Init Job (exchanges, queues, bindings)

3. Proto Dosyaları

  • ingestion_service.proto (eski adı: data_collector.proto)
  • clickhouse_writer.proto (yeni)
  • Message format definitions

4. Dokümantasyon

  • RABBITMQ_INTEGRATION.md - Detaylı entegrasyon kılavuzu
  • MICROSERVICES_NAMES.md - Service isimlendirme standardı
  • MICROSERVICES_ARCHITECTURE.md - Güncellenmiş mimari
  • services/README.md - Microservice genel bakış

5. Scripts

  • scripts/generate_proto.sh - Proto generation script

Deployment Durumu

RabbitMQ

# Deploy komutu
kubectl apply -f deployment/kubernetes-manifests/06-rabbitmq.yaml

# Sonuç
✅ service/rabbitmq created
✅ configmap/rabbitmq-config created
✅ statefulset.apps/rabbitmq created
✅ job.batch/rabbitmq-init created

Status: Pod başlatılıyor...

Management UI


RabbitMQ Configuration

Exchanges (Topic)

  1. flowfish.network_flows - Network flow events
  2. flowfish.dns_queries - DNS query events
  3. flowfish.tcp_connections - TCP connection events

Queues (Durable)

  1. network_flows.clickhouse - Max 1M messages, TTL 24h
  2. dns_queries.clickhouse - Max 1M messages, TTL 24h
  3. tcp_connections.clickhouse - Max 1M messages, TTL 24h

Bindings

flowfish.network_flows → network_flows.clickhouse (#)
flowfish.dns_queries → dns_queries.clickhouse (#)
flowfish.tcp_connections → tcp_connections.clickhouse (#)

Message Flow

1. Ingestion Service (Publisher)

# Connect to RabbitMQ
connection = pika.BlockingConnection(pika.ConnectionParameters('rabbitmq'))
channel = connection.channel()

# Publish message
message = {
    "message_id": str(uuid.uuid4()),
    "event_type": "network_flow",
    "cluster_id": 1,
    "analysis_id": 42,
    "timestamp": datetime.utcnow().isoformat(),
    "data": { ... }
}

channel.basic_publish(
    exchange='flowfish.network_flows',
    routing_key='',
    body=json.dumps(message),
    properties=pika.BasicProperties(
        delivery_mode=2,  # persistent
        content_type='application/json'
    )
)

2. ClickHouse Writer (Consumer)

# Connect to RabbitMQ
connection = pika.BlockingConnection(pika.ConnectionParameters('rabbitmq'))
channel = connection.channel()

# Set QoS (prefetch)
channel.basic_qos(prefetch_count=1000)

# Consume messages
channel.basic_consume(
    queue='network_flows.clickhouse',
    on_message_callback=on_message,
    auto_ack=False
)

def on_message(ch, method, properties, body):
    message = json.loads(body)
    batch.append(message)
    
    if len(batch) >= 1000:
        flush_to_clickhouse(batch)
        ch.basic_ack(delivery_tag=method.delivery_tag, multiple=True)
        batch.clear()

Advantages of RabbitMQ Integration

1. Decoupling

  • Ingestion service ve ClickHouse writer bağımsız
  • Bir service fail olsa diğeri etkilenmiyor

2. Buffering & Backpressure

  • RabbitMQ mesajları queue'da tutar
  • ClickHouse yavaşsa, mesajlar kaybolmaz
  • Ingestion service bloke olmaz

3. Scalability

  • Multiple ingestion workers (1-N)
  • Multiple ClickHouse writers (1-N)
  • Her worker bağımsız scale edilir

4. Reliability

  • Message persistence (disk)
  • Delivery acknowledgment
  • Dead Letter Queue (DLQ)
  • Retry logic

5. Flexibility

  • Aynı veriyi birden fazla consumer okuyabilir
  • Gelecekte: Neo4j writer, Kafka bridge, etc.
  • Real-time streaming to dashboard

Performance Metrics

Expected Throughput

  • Ingestion Service: 10K-50K events/sec (per worker)
  • RabbitMQ: 100K+ messages/sec
  • ClickHouse Writer: 100K-500K rows/sec (bulk insert)

Latency

  • Ingestion → RabbitMQ: < 1ms
  • RabbitMQ → Writer: < 10ms (batch wait)
  • Writer → ClickHouse: < 100ms (bulk insert)
  • End-to-End: < 200ms (event to database)

Next Steps

1. Test RabbitMQ Deployment

# Wait for pod to be ready
kubectl wait --for=condition=ready pod -l app=rabbitmq -n flowfish --timeout=300s

# Port forward Management UI
kubectl port-forward -n flowfish svc/rabbitmq 15672:15672

# Open browser: http://localhost:15672
# Login: flowfish / flowfish-rabbit-2024

# Verify exchanges and queues

2. Implement Ingestion Service

  • Python + gRPC + asyncio + pika
  • Connect to Inspektor Gadget
  • Transform events
  • Publish to RabbitMQ

3. Implement ClickHouse Writer

  • Python + pika + clickhouse-driver
  • Consume from RabbitMQ
  • Batch messages
  • Bulk insert to ClickHouse

4. Integration Test

  • Deploy test Ingestion worker
  • Send sample events
  • Verify messages in RabbitMQ
  • Verify data in ClickHouse

Status Summary

Component Status Next Action
RabbitMQ Architecture Complete -
RabbitMQ Deployment Deployed Wait for pod ready
Proto Files Complete Generate Python code
Documentation Complete -
Ingestion Service 🔴 Not Started Implement
ClickHouse Writer 🔴 Not Started Implement
Integration Test 🔴 Not Started Test end-to-end

Overall Progress: 60% (Architecture & Deployment)
Next Milestone: Microservice Implementation