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)
TermDescription
TopicCategory/channel for messages
PartitionTopics split into partitions (parallel processing)
OffsetSequential ID of messages within a partition
ProducerApplication that sends messages to Kafka
ConsumerApplication that reads messages from Kafka
Consumer GroupGroup of consumers sharing topic partitions
BrokerKafka server node
ZooKeeper/KRaftCluster metadata management

Kafka Use Cases

Use CaseDescription
Log AggregationCollect logs from multiple servers
Event StreamingReal-time event processing
Data IntegrationConnect different systems
Metrics CollectionApplication monitoring
Stream ProcessingProcess events with Kafka Streams / Flink
Website ActivityTrack page views, searches

Kafka vs Traditional Messaging

FeatureTraditional MQ (RabbitMQ)Kafka
Message retentionDeleted after consumptionConfigurable retention (days)
ThroughputLowerVery high (millions/sec)
ReplayNot possibleYes (offset-based)
OrderQueue-basedPer-partition ordering
ScaleLimitedMassive
Use caseTask queuesEvent streaming

Big Data in the Cloud

AWS Big Data Services

ServiceCategoryDescription
Amazon S3StorageObject store; data lake foundation
Amazon RedshiftDWHCloud data warehouse; columnar
Amazon EMRProcessingManaged Hadoop/Spark on AWS
Amazon KinesisStreamingReal-time data streaming
AWS GlueETLServerless data integration
Amazon AthenaQueryQuery S3 data with SQL (serverless)
Amazon QuickSightBICloud BI and visualization

Google Cloud Big Data

ServiceDescription
BigQueryServerless, massive-scale SQL DWH
Cloud DataflowManaged stream + batch processing (Apache Beam)
Cloud Pub/SubManaged Kafka-like messaging
Cloud DataprocManaged Hadoop/Spark
LookerBI and data visualization

Azure Big Data

ServiceDescription
Azure Data Lake StorageScalable data lake
Azure Synapse AnalyticsIntegrated data warehouse + spark
Azure Data FactoryETL/data integration
Azure Event HubsKafka-compatible event streaming
Power BIBI 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.