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
- Start the server: `npm run server`
- Open browser: Navigate to `http://localhost:3001\`
- Configure Kafka: Click "⚙️ Save Config" to set up your Kafka cluster
- Toggle Kafka mode: Use the "🔗 Real Kafka" / "🧪 Mock Kafka" button to switch between real and mock data
Creating a Pipeline
- Add Source Nodes: Enter a topic name and click "+ Source" to create Kafka source nodes
- Add Filter Nodes: Click "+ Filter" to add processing nodes
- Connect Nodes: Drag from output handles to input handles to create connections
- Configure Filters: Click on filter nodes and use "Edit Filter" to set up processing logic
- 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 ```