Distributed Stream Processing with Apache Flink
Introduction: The Real-Time Imperative
- What You Will Learn
- Prerequisites
- How This Book Is Structured
- Version Notes
Chapter 1: The Era of Streaming Data
- From Batch to Streaming: Why Latency Matters
- The Lambda Architecture and Its Pain Points
- From Lambda to Kappa: Unified Processing Models
- Core Challenges of Distributed Stream Processing
- Where Flink Fits in the Streaming Landscape
- Summary
Chapter 2: Stream Processing Fundamentals
- Data Streams, Events, and Event Streams
- Stream vs Batch: Execution Models and Mindset
- Ordering, Latency, and Throughput in Distributed Systems
- Partitioning and Shuffling in Stream Processing
- Fault Tolerance Models: At-Most-Once, At-Least-Once, Exactly-Once
- The Role of Time in Stream Processing
- Summary
Chapter 3: Flink Architecture Overview
- The Flink Runtime: Components and Responsibilities
- Job Submission: From Application Code to Running Job
- The Job Graph: Nodes, Edges, and Stream Transformations
- From Job Graph to Execution Graph
- Task Slots, Parallelism, and Resource Allocation
- The Life Cycle of a Flink Job
- Summary
Chapter 4: Building Your First Flink Application
- Development Environment and Dependencies
- Project Structure for a Flink Application
- Writing a WordCount Streaming Job
- Running Locally and Inspecting Results
- Understanding the Generated Job Graph
- Packaging and Submitting to a Cluster
- Summary
Chapter 5: The DataStream API
- Sources and Sinks: The Boundaries of Your Pipeline
- Core Transformations: Map, FlatMap, Filter, KeyBy
- Reduction and Aggregation Operators
- Process Functions and Fine-Grained Control
- Operator Chaining and Pipelining
- Side Outputs and Branching Streams
- Summary
Chapter 6: Time, Timestamps, and Watermarks
- Processing Time, Event Time, and Ingestion Time
- Why Event Time Matters for Correct Results
- Assigning Timestamps to Streams
- Watermarks: The Heart of Event-Time Processing
- Watermark Strategies and Generators
- Handling Out-of-Order Events
- Summary
Chapter 7: Windowing
- The Concept of Windows
- Time Windows: Tumbling, Sliding, and Session
- Count and Global Windows
- Triggers: Controlling When Windows Fire
- Evictors and Window Functions
- Late Data and Allowed Lateness
- Summary
Chapter 8: State Management Fundamentals
- Why State Changes Everything
- Keyed State vs Operator State
- State Types: Value, List, Map, Reducing, Aggregating
- Accessing and Updating State
- State Backends: Memory, FileSystem, and RocksDB
- Choosing the Right State Strategy
- Summary
Chapter 9: Fault Tolerance: Checkpoints and Savepoints
- The Distributed Snapshot Problem
- Checkpointing in Flink: The Barrier Flow Algorithm
- Configuration and Tuning Checkpoints
- Fault Recovery and State Restoration
- Savepoints: Managed Checkpoints for Operations
- Exactly-Once, At-Least-Once, and At-Most-Once Guarantees
- Summary
Chapter 10: Kafka Integration
- Why Kafka and Flink Are a Natural Pair
- The Flink Kafka Consumer
- The Flink Kafka Producer
- Offset Management and Checkpointing
- Partitioning Semantics and Ordering
- Kafka Integration Patterns and Best Practices
- Summary
Chapter 11: Connectors and External Systems
- Connector Architecture and Interfaces
- Database Connectors: JDBC, Elasticsearch, Redis
- File System Connectors: HDFS, S3, Local
- Change Data Capture: Debezium and Flink CDC
- Building Custom Sources and Sinks
- Choosing Connectors for Your Use Case
- Summary
Chapter 12: The Flink Table API and SQL
- The Table API: Declarative Stream Processing
- Table Environments and Execution Mode
- Defining Tables and Schemas
- SQL for Streaming: Queries and Semantics
- Bridging DataStream and Table APIs
- When to Use SQL vs DataStream
- Summary
Chapter 13: Stream Joins and Enrichment
- Joining Event Streams: The Challenge
- Stream-Stream Joins: Interval and Window Joins
- Stream-Table Joins: Enriching with Reference Data
- Temporal Table Joins and Schema Evolution
- Broadcast State for Lookup and Enrichment
- Handling Asynchrony and Skewed Data
- Summary
Chapter 14: Advanced Data Processing Patterns
- Deduplication and Idempotency
- Anomaly Detection on Streams
- Sessionization and User Activity Tracking
- Multi-Dimensional Aggregations
- Multi-Stage Pipelines and Topology Design
- Real-Time Feature Computation for ML
- Summary
Chapter 15: Serialization and Data Formats
- Flink’s Type System
- Serialization: Kryo, Java, Generic Types
- Optimizing Serialization Performance
- Avro, Protobuf, and Schema Management
- Schema Evolution in Streaming
- Binary Formats and Network Efficiency
- Summary
Chapter 16: Parallelism, Backpressure, and Scaling
- Parallelism Model and Task Allocation
- Understanding and Detecting Backpressure
- Rescaling Jobs Without Losing State
- Data Skew and Its Impact
- Slot Sharing and Co-Localization
- Elastic Scaling Patterns
- Summary
Chapter 17: Memory Management and Resource Tuning
- Flink’s Memory Layout
- Task, Network, and Managed Memory
- Tuning TaskManager and JobManager Resources
- JVM Garbage Collection and Flink
- Monitoring Memory Pressure
- Memory-Related Failure Modes
- Summary
Chapter 18: Performance Tuning and Optimization
- Latency vs Throughput Trade-Offs
- Optimizing State Access and Size
- Checkpoint Optimization Strategies
- Network Shuffle Optimization
- Operator-Level Optimization Techniques
- Profiling and Diagnosing Performance Problems
- Summary
Chapter 19: Deployment and Operations
- Deployment Modes: Standalone, YARN, Kubernetes
- Cluster Modes: Session, Application, Per-Job
- High Availability Configuration
- Resource Management and Cluster Sizing
- Configuration Management
- Production Deployment Checklist
- Summary
Chapter 20: Observability, Monitoring, and Debugging
- The Flink Web Dashboard
- Built-In Metrics and Custom Metrics
- Integrating with Prometheus and Grafana
- Logging Strategies for Distributed Jobs
- Debugging Techniques and Tools
- Incident Response and Troubleshooting
- Summary
Chapter 21: Security, Testing, and Reliability
- Authentication and Authorization
- Encryption in Transit and at Rest
- Network Security and Kerberos
- Testing Flink Applications
- CI/CD Pipelines for Flink Jobs
- Production Reliability Patterns
- Summary
Chapter 22: Production Architecture Patterns and Anti-Patterns
- End-to-End Real-Time Analytics Architecture
- Fraud Detection Architecture
- CDC-Based Event-Driven Architecture
- Multi-Tenant Flink Platforms
- Common Anti-Patterns and How to Avoid Them
- Version Upgrades and Long-Term Maintenance
- Summary