Big data & Spark
PySpark jobs, partitioning, batch processing, and large-dataset transform patterns.
How to Simulate a MapReduce Mock with Combine Phase in Python
Simulates a MapReduce pipeline with a combiner that aggregates local counts per reducer to reduce network and compute overhead.
from collections import defaultdict
def map_phase(lines):
intermediate = defaultdict(list)
for line in lines:
for word in line.strip().lower().split():
intermediate[word].append(1)
return dict(intermediate)
def combine_phase(intermediate, num_reducers=3):
combined = defaultdict(li…
How to Truncate Lineage Back to a Checkpoint in Python
Walks a linked list of lineage nodes upward to find the nearest checkpoint and returns that node, truncating the lineage.
class LineageNode:
def __init__(self, name, parent=None, checkpoint=None):
self.name = name
self.parent = parent
self.checkpoint = checkpoint
def truncate_at_checkpoint(self):
"""Truncate lineage back to the last checkpoint."""
current = self
while current.check…
How to Use Broadcast Variables as Read-Only in PySpark (Mock Example)
Share a lookup dict across Spark executors with a broadcast variable and verify its read-only behavior in a local mock.
from pyspark import SparkContext, SparkConf
def main():
conf = SparkConf().setAppName("BroadcastMock").setMaster("local[2]")
sc = SparkContext(conf=conf)
lookup = {"a": 1, "b": 2, "c": 3}
broadcast_lookup = sc.broadcast(lookup)
data = ["a", "b", "c", "a", "unknown"]
rdd = sc.parallel…
How to implement a tumbling window aggregation in Python
Build a mock tumbling window aggregator in Python that groups streaming events into fixed time intervals and computes count, sum, and average per window.
import time
from collections import deque
class TumblingWindow:
def __init__(self, duration_seconds):
self.duration = duration_seconds
self.buffer = deque()
self.window_start = None
def add(self, item):
current_time = time.time()
if self.window_start is None:
…
How to select specific columns in Python with SQLite
A reusable function that connects to a SQLite database and returns only the requested columns from a given table.
import sqlite3
def select_pruned_columns(db_path, table, columns):
with sqlite3.connect(db_path) as conn:
cursor = conn.cursor()
col_list = ", ".join(columns)
query = f"SELECT {col_list} FROM {table}"
return cursor.execute(query).fetchall()
if __name__ == "__main__":
conn = sq…
How to use foreachBatch with a mock sink in PySpark
Demonstrates using Spark Structured Streaming's foreachBatch sink to capture and verify streaming batches by writing them into a custom mock sink object.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit
class MockSink:
def __init__(self):
self.batches = []
def write_batch(self, batch_df, batch_id):
# Collect batch data as list of dicts for verification
records = batch_df.collect()
self.batches…
Hudi Upsert Mock Copy on Write in Python
Simulates Apache Hudi's Copy-on-Write upsert behavior by merging update records into a deep copy of base records, replacing matches or appending new ones.
import copy
from typing import Dict, List, Any
def upsert_copy_on_write(base_records: List[Dict[str, Any]], updates: List[Dict[str, Any]], key_field: str = "id") -> List[Dict[str, Any]]:
"""Simulate Hudi Copy-on-Write upsert: merge updates into a copy of base records."""
result = copy.deepcopy(base_records)
…
HyperLogLog Cardinality Estimation in Python
A small HyperLogLog implementation using MD5 hashing and 256 registers to estimate the number of unique items in a large stream with fixed memory.
import hashlib
import math
class HyperLogLog:
def __init__(self, b=8):
self.b = b
self.m = 1 << b
self.registers = [0] * self.m
self.alpha = 0.7213 / (1 + 1.079 / self.m)
def add(self, item):
h = int(hashlib.md5(str(item).encode()).hexdigest(), 16)
idx = h & (s…
Mock Predicate Pushdown in Python for Big Data Queries
Simulate predicate pushdown by applying filters at the storage layer before materializing rows, showing how big data engines optimize queries.
class Query:
def __init__(self, table, rows):
self.table = table
self.rows = rows
def filter(self, predicate):
return Query(
self.table,
[row for row in self.rows if all(predicate(row) for predicate in predicate)]
)
def filter_pushdown(self, predica…
Modeling a Hive Metastore Table Schema in Python
A dataclass that mimics a Hive metastore table schema—columns, partition keys, storage format, and location—with helper methods for description and mutation.
from dataclasses import dataclass, field
from typing import Dict, List, Optional
@dataclass
class HiveTable:
"""Simple mock of a Hive metastore table schema."""
name: str
database: str = "default"
columns: List[Dict[str, str]] = field(default_factory=list)
partition_keys: List[Dict[str, str]] = f…
Partition Data by Hash Key Mod N in Python
Returns a partition index for a string key by hashing it with MD5 and taking modulo N, then groups sample keys into partitions.
import hashlib
def partition_key(key: str, num_partitions: int) -> int:
"""Return partition index for key using MD5 hash mod N."""
digest = hashlib.md5(key.encode()).hexdigest()
return int(digest, 16) % num_partitions
if __name__ == "__main__":
keys = ["alice", "bob", "carol", "dave", "eve"]
nu…
Session window gap mock in Python
Group sorted timestamps into sessions where any gap between consecutive events exceeds a threshold starts a new session.
from datetime import datetime, timedelta
def session_windows(timestamps, gap_seconds=300):
"""Group timestamps into sessions where gaps > gap_seconds start new sessions."""
if not timestamps:
return []
# Sort timestamps chronologically to ensure correct windowing
timestamps = sorted(timestam…
Sliding Window Streaming Mock in Python
A simple Python class that maintains a sliding window of recent streaming values and computes the running average.
import time
import random
class StreamingMock:
"""Produces a stream of numbers using a sliding window."""
def __init__(self, window_size=5):
self.window = []
self.window_size = window_size
def push(self, value):
"""Add a value, sliding the window forward."""
s…
Z-Order Optimization in Python
A mock concept demonstrating z-order layout optimization by reassigning z-indices based on areas size.
class ZOrderLayout:
"""
Minimal mock for z-order layout optimization using a stacking score.
Elements overlap; higher z_index is drawn on top.
"""
def __init__(self):
self.elements = []
def add_element(self, name, area, z_index):
self.elements.append({"name": name, "area": area…
Browse by section
Each section groups closely related Python snippets.
Big data & Spark — Python code examples
What you will find here
This page collects big data & spark snippets — short, copy-ready Python you can paste into our free online IDE and run without installing anything. Each sample includes a plain-English explanation and the full source code.
Samples vs tutorials and challenges
Samples are quick reference — one concept per page. For step-by-step teaching, use our Python tutorials. To test yourself, try quizzes or coding challenges. Clean up style with the Python formatter.