Levi DeHaan

Kafka Visual Data Pipeline

Comprehensive full-stack application for creating, managing, and monitoring Kafka data pipelines through an intuitive visual interface with real-time processing, AI filters, and drag-and-drop pipeline building.

Status: completed · 2025-04-01

Overview

A comprehensive full-stack application for creating, managing, and monitoring Kafka data pipelines through an intuitive visual interface. Built with React, ReactFlow, Node.js, and real Kafka integration, featuring drag-and-drop pipeline building, real-time data flow visualization, advanced filtering with AI processing, and live monitoring.

Technologies

React 18, TypeScript, Node.js, Express.js, KafkaJS, ReactFlow, Tailwind CSS, WebSocket, Kafka Cluster Integration, AI Services (Gemini, OpenAI, Claude), Custom JavaScript Filters, Real-time Data Processing, JSON Configuration, Pipeline Execution Engine, Error Handling & Monitoring, Responsive UI, API Design, WebSocket Communication, Drag-and-Drop Interface, Data Flow Animation

Real-time Data Processing
WebSocket-based
AI Filter Integration
Multi-provider
Pipeline Nodes
Drag-and-drop
Kafka Topics
Dynamic Management

Kafka Visual Data Pipeline

A comprehensive full-stack application for creating, managing, and monitoring Kafka data pipelines through an intuitive visual interface. Built with React, ReactFlow, Node.js, and real Kafka integration.

🚀 Features

Visual Pipeline Builder

  • Drag-and-drop interface using ReactFlow for creating data pipelines
  • Real-time data flow visualization with animated data items
  • Node-based architecture with source, filter, and output nodes
  • Interactive connection management between pipeline components

Real Kafka Integration

  • Live Kafka cluster connectivity (default: 192.168.8.161:9092)
  • Topic management (create, delete, list topics)
  • Real-time message consumption and production
  • WebSocket-based real-time updates

Advanced Filtering System

  • Field Matching Filters: Filter data based on field values and conditions
  • AI Processing Filters: Integrate with AI services (Gemini, OpenAI, Claude) for intelligent data processing
  • Custom Code Filters: Execute custom JavaScript code for complex data transformations
  • Filter chaining and complex pipeline logic

Configuration Management

  • Persistent configuration with JSON-based storage
  • Kafka cluster settings (brokers, timeouts, retry logic)
  • Pipeline configuration (auto-start, error handling, concurrency)
  • AI provider setup and API key management
  • UI customization (themes, animations, refresh intervals)

Real-time Monitoring

  • Live pipeline logs with node-specific and global views
  • Data flow animations showing message processing
  • Error tracking and status monitoring
  • WebSocket-based real-time updates

🏗️ Architecture

Frontend (React + TypeScript)

  • React 18 with TypeScript for type safety
  • ReactFlow for visual pipeline creation
  • Tailwind CSS for modern, responsive UI
  • WebSocket client for real-time communication
  • Service layer for API communication

Backend (Node.js + Express)

  • Express.js server with comprehensive API
  • KafkaJS for real Kafka cluster integration
  • WebSocket server for real-time updates
  • Configuration management with persistent storage
  • Pipeline execution engine with filter processing

Services Architecture

``` Frontend (React) ↓ HTTP/WebSocket Backend Server (Express) ↓ KafkaJS Kafka Cluster (192.168.8.161:9092) ```

📦 Installation

Prerequisites

  • Node.js 18+ and npm
  • Access to a Kafka cluster (default: 192.168.8.161:9092)
  • Modern web browser with WebSocket support

Setup

```bash

Clone the repository

git clone cd kafka-visual-data-pipeline

Install dependencies

npm install

Build the frontend

npm run build

Start the backend server (includes frontend serving)

npm run server ```

Development Mode

```bash

Start backend in development mode

npm run dev:backend

In another terminal, start frontend development server

npm run dev ```

🚀 Usage

Starting the Application

  1. Start the server: `npm run server`
  2. Open browser: Navigate to `http://localhost:3001\`
  3. Configure Kafka: Click "⚙️ Save Config" to set up your Kafka cluster
  4. Toggle Kafka mode: Use the "🔗 Real Kafka" / "🧪 Mock Kafka" button to switch between real and mock data

Creating a Pipeline

  1. Add Source Nodes: Enter a topic name and click "+ Source" to create Kafka source nodes
  2. Add Filter Nodes: Click "+ Filter" to add processing nodes
  3. Connect Nodes: Drag from output handles to input handles to create connections
  4. Configure Filters: Click on filter nodes and use "Edit Filter" to set up processing logic
  5. Monitor Data Flow: Watch real-time data items flow through your pipeline

Filter Types

Field Matching Filter

```javascript // Example: Filter messages where status equals "active" { field: "status", operator: "equals", value: "active" } ```

AI Processing Filter

```javascript // Example: Sentiment analysis with Gemini { provider: "gemini", prompt: "Analyze the sentiment of this message and return positive/negative/neutral", outputField: "sentiment" } ```

Custom Code Filter

```javascript // Example: Custom transformation function processData(data) { return { ...data, processed: true, timestamp: new Date().toISOString(), score: data.value * 2 }; } ```

⚙️ Configuration

Kafka Settings

  • Client ID: Unique identifier for your Kafka client
  • Brokers: List of Kafka broker addresses
  • Connection Timeout: Timeout for initial connections
  • Request Timeout: Timeout for individual requests
  • Retry Logic: Initial retry time and maximum retries

Pipeline Settings

  • Auto-start: Automatically start pipeline on server startup
  • Max Concurrent Messages: Limit concurrent message processing
  • Error Handling: Choose between log, stop, or ignore error strategies
  • Retry Failed Messages: Enable/disable message retry logic

AI Integration

  • Provider: Choose between Gemini, OpenAI, or Anthropic
  • API Key: Your AI service API key
  • Model: Specific model to use (e.g., gemini-2.5-flash-preview-04-17)
  • Timeout: Request timeout for AI processing

UI Customization

  • Theme: Light or dark mode
  • Auto-refresh: Enable automatic data refresh
  • Refresh Interval: How often to refresh data
  • Show Animations: Enable/disable data flow animations
  • Max Data Items: Limit visible data items for performance

🔧 API Endpoints

Health & Status

  • `GET /api/health` - Server health check
  • `GET /api/pipeline` - Get pipeline status

Configuration

  • `GET /api/config` - Get current configuration
  • `POST /api/config` - Save configuration

Kafka Management

  • `GET /api/kafka/topics` - List all topics
  • `POST /api/kafka/topics` - Create new topic
  • `DELETE /api/kafka/topics/:name` - Delete topic
  • `POST /api/kafka/publish` - Publish message to topic

Pipeline Control

  • `POST /api/pipeline/start` - Start pipeline processing
  • `POST /api/pipeline/stop` - Stop pipeline processing
  • `POST /api/filters/execute` - Execute filter on data

AI Processing

  • `POST /api/ai/process` - Process data with AI

🔌 WebSocket Events

Client → Server

  • `subscribe_topic` - Subscribe to topic messages
  • `unsubscribe_topic` - Unsubscribe from topic
  • `get_pipeline_status` - Request pipeline status

Server → Client

  • `topic_message` - New message from subscribed topic
  • `data_received` - Data received in pipeline
  • `filter_processed` - Filter processing completed
  • `data_rejected` - Data rejected by filter
  • `data_output` - Data sent to output
  • `pipeline_status_changed` - Pipeline status update
  • `error` - Error occurred

📁 Project Structure

``` kafka-visual-data-pipeline/ ├── backend/ # Backend server code │ ├── services/ │ │ ├── kafkaService.js # Real Kafka integration │ │ ├── configManager.js # Configuration management │ │ └── pipelineManager.js # Pipeline execution engine │ └── server.js # Main Express server ├── components/ # React components │ ├── SourceNodeComponent.tsx # Kafka source nodes │ ├── FilterNodeComponent.tsx # Filter processing nodes │ ├── OutputNodeComponent.tsx # Output nodes │ ├── FilterEditorModal.tsx # Filter configuration UI │ ├── ConfigModal.tsx # System configuration UI │ └── LogDock.tsx # Real-time logging panel ├── services/ # Frontend services │ ├── mockKafkaService.ts # Mock Kafka for development │ ├── realKafkaService.ts # Real Kafka API client │ └── aiService.ts # AI integration service ├── App.tsx # Main application component ├── types.ts # TypeScript type definitions ├── constants.ts # Application constants └── package.json # Dependencies and scripts ```

🛠️ Development

Available Scripts

  • `npm run dev` - Start frontend development server
  • `npm run dev:backend` - Start backend in development mode
  • `npm run build` - Build frontend for production
  • `npm run server` - Start production server
  • `npm start` - Start full-stack application

Environment Variables

```bash PORT=3001 # Server port KAFKA_BROKERS=192.168.8.161:9092 # Kafka broker addresses ```

🔍 Troubleshooting

Common Issues

Kafka Connection Failed

  • Verify Kafka cluster is running and accessible
  • Check broker addresses in configuration
  • Ensure network connectivity to Kafka cluster

WebSocket Connection Issues

  • Check if port 3001 is available
  • Verify firewall settings
  • Try refreshing the browser

AI Processing Errors

  • Verify API key is correctly configured
  • Check AI provider service status
  • Ensure sufficient API quota/credits

Debug Mode

Enable debug logging by setting environment variable: ```bash DEBUG=kafka-visual-pipeline npm run server ```