Big Data Technologies
Spark Structured Streaming, Kafka, Connect and Real-Time Pipelines
PGCP-BDA
structured streaming
Spark Structured Streaming models incoming data as an incrementally updated table and executes supported DataFrame operations as a fault-tolerant streaming.
unbounded table model
Structured Streaming treats an input stream as an unbounded table and maintains incremental results as new rows are appended.
event time and processing time
Event time is when the source event occurred, while processing time is when the engine handled it
watermark
A watermark estimates how far event time has progressed and bounds how long state is retained for late records.
output mode
A streaming output mode chooses whether each trigger emits only new rows, changed result rows or the complete result table where supported.
checkpoint
A streaming checkpoint persists offsets, progress and state-store metadata so a query can resume consistently after driver failure.
Kafka broker topic partition
A Kafka broker stores replicated topic partitions, each an ordered append-only log whose offsets identify record positions.
producer and consumer group
A Kafka producer chooses a topic partition for records, while consumers in one group divide those partitions so each partition has one active group member.
Kafka offset and delivery semantics
An offset identifies a record position in a Kafka partition; commit and retry behavior determines at-most-once, at-least-once or exactly-once processing.
Spark Kafka integration
Spark’s Kafka connector maps Kafka records into streaming DataFrame rows and coordinates source offsets with checkpointed query progress.
Apache Kafka
What is Kafka?
Apache Kafka is a distributed event streaming platform for:
- High-throughput, low-latency message streaming
- Real-time data pipelines
- Decoupling of data producers and consumers
Kafka Core Concepts
PRODUCER (writes data) KAFKA CLUSTER CONSUMER (reads data)
+-----------+ +----------+ +----------+
| App | produces | Topic | reads | Analytics|
| Database |──────────────| (logs) |────────| ML model |
| IoT device| | Partition| | Dashboard|
| | | 1, 2, 3 | | |
+-----------+ +----------+ +----------+
Brokers (servers)
| Term | Description |
|---|---|
| Topic | Category/channel for messages |
| Partition | Topics split into partitions (parallel processing) |
| Offset | Sequential ID of messages within a partition |
| Producer | Application that sends messages to Kafka |
| Consumer | Application that reads messages from Kafka |
| Consumer Group | Group of consumers sharing topic partitions |
| Broker | Kafka server node |
| ZooKeeper/KRaft | Cluster metadata management |
Kafka Use Cases
| Use Case | Description |
|---|---|
| Log Aggregation | Collect logs from multiple servers |
| Event Streaming | Real-time event processing |
| Data Integration | Connect different systems |
| Metrics Collection | Application monitoring |
| Stream Processing | Process events with Kafka Streams / Flink |
| Website Activity | Track page views, searches |
Kafka vs Traditional Messaging
| Feature | Traditional MQ (RabbitMQ) | Kafka |
|---|---|---|
| Message retention | Deleted after consumption | Configurable retention (days) |
| Throughput | Lower | Very high (millions/sec) |
| Replay | Not possible | Yes (offset-based) |
| Order | Queue-based | Per-partition ordering |
| Scale | Limited | Massive |
| Use case | Task queues | Event streaming |
Big Data in the Cloud
AWS Big Data Services
| Service | Category | Description |
|---|---|---|
| Amazon S3 | Storage | Object store; data lake foundation |
| Amazon Redshift | DWH | Cloud data warehouse; columnar |
| Amazon EMR | Processing | Managed Hadoop/Spark on AWS |
| Amazon Kinesis | Streaming | Real-time data streaming |
| AWS Glue | ETL | Serverless data integration |
| Amazon Athena | Query | Query S3 data with SQL (serverless) |
| Amazon QuickSight | BI | Cloud BI and visualization |
Google Cloud Big Data
| Service | Description |
|---|---|
| BigQuery | Serverless, massive-scale SQL DWH |
| Cloud Dataflow | Managed stream + batch processing (Apache Beam) |
| Cloud Pub/Sub | Managed Kafka-like messaging |
| Cloud Dataproc | Managed Hadoop/Spark |
| Looker | BI and data visualization |
Azure Big Data
| Service | Description |
|---|---|
| Azure Data Lake Storage | Scalable data lake |
| Azure Synapse Analytics | Integrated data warehouse + spark |
| Azure Data Factory | ETL/data integration |
| Azure Event Hubs | Kafka-compatible event streaming |
| Power BI | BI visualization |
Continue learning
Related notes
Put this topic into timed practice
Open mock tests when you want full-exam pacing, or keep drilling in practice mode.