\

Chapter 3: Data Engineering Fundamentals

25 min read

As the book says, “The rise of ML in recent years is tightly coupled with the rise of big data.” At FAANG, this is our daily reality. Our ML systems, from recommendation engines to fraud detection, are built on vast, complex data landscapes. If you don’t have a solid grasp of data engineering, you’re going to struggle, no matter how fancy your model architecture is.

This chapter is dense with terminology and concepts that might seem overwhelming if you’re new to large-scale data systems. But don’t worry, we’ll break it down. The goal here is to give you a “steady piece of land to stand on.”

We’ll cover:

  • Data Sources: Where does our data even come from?
  • Data Formats: How is it stored? What are the trade-offs?
  • Data Models: How is it structured and represented?
  • Data Storage Engines & Processing: Databases, transactional vs. analytical workloads, ETL.
  • Modes of Dataflow: How does data move between different parts of our systems?
  • Batch vs. Stream Processing: Two fundamental paradigms for handling data.

This chapter is critical for ML system design interviews. Understanding data is non-negotiable. Let’s get started!

Page 49 (Chapter Introduction): The Data Deluge

The intro page sets the scene:

  • ML’s growth is linked to big data.
  • Large data systems are complex, even without ML. Full of acronyms, evolving standards, diverse tools. It can feel like every company does it differently.
  • This chapter provides the basics of data engineering.

The roadmap mentioned:

  • Sources of data.
  • Formats for storage.
  • Structure of data (data models).
  • Databases (storage engines) for transactional and analytical processing.
  • Passing data across multiple processes and services.

This chapter really emphasizes that data engineering is a distinct discipline that ML practitioners must understand.

Pages 50-52: Data Sources – Where It All Begins

An ML system ingests data from various sources. Understanding these sources helps use the data more efficiently and anticipate challenges.

User Input Data (Page 50)

Data explicitly provided by users.

  • Examples: Text typed into a search bar, images/videos uploaded, form submissions.
  • Challenges:
    • Malformatted: “If it’s even remotely possible for users to input wrong data, they are going to do it.” This is a golden rule! Text too long/short, text in numerical fields, wrong file formats.
    • Requires heavy-duty validation and processing.
    • User Patience is Low: Users expect immediate results from their input. This implies a need for fast processing.

FAANG Perspective: Input validation is a huge deal. For example, search queries need sanitization, image uploads need format checks and virus scans. Latency is critical – a slow search box is a bad user experience.

System-Generated Data (Pages 50-51)

Data generated by your system’s components.

  • Examples: Logs (memory usage, services called, job results), model predictions.
  • Purpose: Visibility into system health, debugging, improving the application. Essential “when something is on fire.”
  • Characteristics:
    • Less likely to be malformatted than user input.
    • Processing doesn’t always need to be immediate (e.g., hourly/daily log processing is often fine).
    • However, you might want faster processing for “interesting” events (footnote 1: “interesting” often means “catastrophic” – like a system crash or a runaway cloud bill!). This is where real-time alerting on logs comes in.
  • Challenges with Logs:
    • Volume: “Log everything you can” is common practice for debugging ML systems, leading to massive log volumes.
    • Signal vs. Noise: Hard to find useful information. Services like Logstash, Datadog, Logz.io (often using ML themselves) help process and analyze logs.
    • Storage Cost: Store logs only as long as useful. Low-access, cheaper storage (e.g., AWS S3 Glacier vs. S3 Standard – footnote 2 notes a 5x cost difference for much higher retrieval latency) can be used for older logs.

FAANG Perspective: We generate petabytes of logs daily. Effective log management, aggregation, and analysis are critical for SREs (Site Reliability Engineers) and ML engineers alike to debug production issues.

User Behavior Data (Page 51)

System-generated data specifically recording user actions.

  • Examples: Clicks, suggestions chosen, scrolling, zooming, ignoring pop-ups, time spent on page.
  • Crucial for ML: This is the raw material for training many models (recommenders, personalization, engagement prediction).
  • Privacy Concerns: Even if system-generated, it’s considered user data and subject to privacy regulations (e.g., GDPR, CCPA). Footnote 3 has a great anecdote: an ML engineer says his team only uses browsing/purchase history, not “personal data” like age/location. The author rightly points out that browsing/purchase history is extremely personal!

FAANG Perspective: This is gold. But handling it ethically and in compliance with regulations is paramount. Data anonymization and aggregation are key, but true anonymization is very hard.

Internal Databases (Page 52)

Generated by various services and enterprise applications within the company.

  • Examples: Inventory, customer relationship management (CRM), user accounts.
  • Usage in ML:
    • Directly by models (e.g., a model might need current user subscription status).
    • By components of an ML system (e.g., Amazon search: ML model detects intent for “frozen,” then system checks internal inventory database for “Frozen” movie vs. “frozen foods” availability before ranking).

FAANG Perspective: These are often the “source of truth” for many entities. Integrating them into ML pipelines reliably and efficiently is a common data engineering task.

Third-Party Data (Page 52)

The “wonderfully weird world.”

  • First-party data: What your company collects about its own users/customers. (This is the best, you own it, you know its provenance).
  • Second-party data: Data from another company on their customers, which they make available to you (usually for a fee).
  • Third-party data: Companies collect data about the public (not their direct customers) and sell it.
  • Collection Mechanisms: Historically, unique advertiser IDs on phones (IDFA on iOS, AAID on Android) made it easy to aggregate activity across apps, websites, check-ins. This data is (hopefully) anonymized.
  • Types of Data: Social media activity, purchase history, web browsing, car rentals, political leaning, demographics (e.g., “men, age 25-34, tech workers, Bay Area”).
  • Use Cases: Inferring correlations (people liking Brand A also like Brand B), helpful for recommenders. Usually sold pre-cleaned and processed.
  • Privacy Pushback:
    • Apple’s IDFA opt-in (early 2021) significantly reduced third-party data on iPhones, forcing companies to focus more on first-party data (footnote 4).
    • Advertisers seek workarounds (e.g., CAID device fingerprinting in China - footnote 5).

FAANG Perspective: While first-party data is king, third-party data can be useful for cold-start problems or enriching user profiles. However, the privacy and ethical implications are significant, and reliance on it is decreasing due to regulations and platform changes.

Understanding your data sources tells you about its likely quality, freshness, volume, and any constraints (privacy, cost, processing needs).

Pages 53-57: Data Formats – How Data is Represented for Storage and Transmission

Once you have data, you need to store it (“persist” it). The format matters for cost, access speed, and ease of use. Key questions to consider:

  • Storing multimodal data (images + text)?
  • Cheap and fast access?
  • Storing complex models for cross-hardware compatibility?

Data Serialization: Converting a data structure/object state into a format that can be stored/transmitted and later reconstructed.

Table 3-1: Common Data Formats

FormatBinary/TextHuman-readableExample Use Cases
JSONTextYesEverywhere
CSVTextYesEverywhere
ParquetBinaryNoHadoop, Amazon Redshift
AvroBinary primaryNoHadoop
ProtobufBinary primaryNoGoogle, TensorFlow (TFRecord)
PickleBinaryNoPython, PyTorch serialization (models)

Key characteristics to consider: human readability, access patterns (how data is read/written - footnote 6), text vs. binary (impacts file size).

Let’s look at a few in detail:

JSON (JavaScript Object Notation) (Page 54)

  • Ubiquitous, language-independent, human-readable.
  • Key-value pair paradigm, handles different levels of structuredness.
    • Structured: {"firstName": "Boatie", "address": {"city": "Port Royal"}}
    • Unstructured blob: {"text": "Boatie McBoatFace, aged 12, is vibing..."}
  • Pain points:
    • Schema evolution is painful once committed. Changing a schema in existing JSON files is hard.
    • Text files = take up a lot of space (more on this later).

FAANG Perspective: JSON is incredibly common for API responses, configs, and semi-structured data logging. Its human readability is a big plus for debugging.

Row-Major Versus Column-Major Format (Pages 54-55, Figure 3-1)

  • CSV (Comma-Separated Values): Row-major. Consecutive elements in a row are stored together in memory.
    Example1_Feat1, Example1_Feat2, ... Example1_FeatN
    Example2_Feat1, Example2_Feat2, ... Example2_FeatN
    
  • Parquet: Column-major (Columnar). Consecutive elements in a column are stored together.
    Example1_Feat1, Example2_Feat1, ... ExampleM_Feat1
    Example1_Feat2, Example2_Feat2, ... ExampleM_Feat2
    
  • Performance implications (due to sequential data access being efficient):
    • Row-major (CSV): Better for accessing entire rows/examples (e.g., “get all examples from today”). Faster writes when adding new individual examples.
    • Column-major (Parquet): Better for accessing specific columns/features (e.g., “get timestamps for all examples”). More flexible column-based reads, especially with many features. If you have 1000 features but only need 4 (time, location, distance, price), Parquet lets you read just those 4 columns directly. With CSV, you’d often read all 1000 and then filter.
  • Overall: Row-major for write-heavy workloads or row-based reads. Column-major for read-heavy, analytical workloads needing specific columns.

FAANG Perspective: Columnar formats like Parquet (and ORC) are the standard for data warehouses and data lakes (e.g., in S3, GCS) because analytical queries usually operate on subsets of columns. This leads to huge I/O savings and faster query performance. Compression is also more effective on columnar data.

NumPy Versus pandas (Page 56, Figure 3-2): A Common Gotcha!

  • Many don’t realize: pandas DataFrame is built around a columnar format (inspired by R’s data frame).
  • NumPy ndarrays are row-major by default (though configurable).
  • People coming from NumPy often treat DataFrames like ndarrays, accessing by row, and find it slow.
  • Figure 3-2 Performance:
    • Iterating pandas DataFrame by column: 0.07 seconds.
    • Iterating pandas DataFrame by row (df.iloc[i]): 2.41 seconds (MUCH SLOWER!).
    • Converting to NumPy array (df.to_numpy()) and iterating by row: 0.019 seconds (FAST!).
    • Iterating NumPy array by column: 0.005 seconds (FASTEST, as expected for columnar data if NumPy array was C-contiguous/columnar, but it’s likely F-contiguous/row-major here, so this illustrates pandas’ overhead).

Self-Correction/Teaching Point: The code snippet df_np[:, j] in Figure 3-2 iterates through a NumPy array column by column. If df_np is row-major (NumPy default), this is non-contiguous access, which should be slower than row-wise access. The 0.005s vs 0.019s suggests the test DataFrame might have few rows and many columns, or there’s something subtle about to_numpy() memory layout or caching effects. The main takeaway is pandas row iteration is slow due to its internal structure and overhead. Use vectorized operations in pandas, or convert to NumPy for row-wise loops if absolutely necessary. The author’s “Just pandas Things” GitHub repo (footnote 7) is a good resource for these quirks.

Text Versus Binary Format (Page 57, Figure 3-3)

  • CSV, JSON are text files (plain text, human-readable).
  • Parquet is a binary file (0s and 1s, machine-readable).
  • Binary files are more compact:
    • Example: Storing the number 1000000.
      • Text file: 7 characters, 7 bytes (if 1 byte/char).
      • Binary file (int32): 32 bits = 4 bytes. (Significant saving!)
  • Figure 3-3 Illustration (interviews.csv):
    • CSV (text): 17,654 rows, 10 columns. File size: 14 MB.
    • Parquet (binary): Same data. File size: 6 MB. (Over 2x smaller).
  • AWS recommends Parquet: “up to 2x faster to unload and consumes up to 6x less storage in Amazon S3, compared to text formats” (footnote 8). This is due to efficient encoding and compression schemes that work well with columnar data.

FAANG Perspective: For large datasets, binary columnar formats (Parquet, ORC) are almost always preferred over text formats like CSV/JSON for storage in data lakes due to space savings, query performance, and schema evolution support. Text formats are fine for smaller files, human inspection, or system interchange where readability is key.

Pages 58-66: Data Models – Structuring Your Data

Data models describe how data is represented and the relationships between data elements. This choice affects system build and the problems you can solve.

Relational Model (Pages 59-61)

  • Invented by Edgar F. Codd (1970), still dominant.
  • Data organized into relations (tables), each a set of tuples (rows).
  • Key Property (Figure 3-4): Relations are unordered (rows and columns can be shuffled, it’s still the same relation). Stored in formats like CSV/Parquet.
  • Normalization: Reducing redundancy, improving integrity.
    • Example (Tables 3-2, 3-3, 3-4): Book data.
      • Initial Book relation (Table 3-2): Title, Author, Format, Publisher, Country, Price. Duplicates publisher info (Banana Press, UK) for different books/formats. If “Banana Press” changes to “Pineapple Press”, multiple rows need updates.
      • Normalized:
        • Book relation (Table 3-3): Title, Author, Format, Publisher_ID, Price.
        • Publisher relation (Table 3-4): Publisher_ID, Publisher, Country.
      • Now, if publisher name changes, only one row in Publisher table needs update. Standardizes spelling, easier to translate values.
  • Downside of Normalization: Data spread across tables. Retrieving full info requires joins, which can be expensive for large tables.
  • Relational Databases & SQL (Structured Query Language):
    • SQL is the most popular query language.
    • Declarative: You specify what data you want (pattern, conditions, transformations like join, sort, group, aggregate), not how to get it.
    • Imperative (like Python): You specify the steps.
    • Database system has a query optimizer to figure out the execution plan (break query, methods for each part, order of execution). This is hard! ML is even being used to improve query optimizers (footnote 14, Neo).
    • SQL is Turing-complete (with additions), but complex queries can be “nightmarish” (footnote 12, 700-line SQL query).

FAANG Perspective: Relational databases (PostgreSQL, MySQL, Spanner, Aurora) are workhorses for many transactional systems and structured data stores. SQL is a fundamental skill. Understanding query plans (EXPLAIN) is key for performance tuning.

Aside: From Declarative Data Systems to Declarative ML Systems (Page 62)

  • Inspired by SQL’s success, “Declarative ML” aims to abstract away model construction/tuning.
  • User declares feature schema and task; system finds the best model.
  • Examples: Ludwig (Uber), H2O AutoML.
    • Ludwig: User can specify model structure (layers, units) on top of schema.
    • H2O AutoML: No need to specify structure/hyperparameters; it experiments and picks best. Example code shows simple API: aml.train(x=x, y=y, training_frame=train).
  • Limitation: Abstracts away model development (often the easier part now with commoditized models). Hard parts remain: feature engineering, data processing, evaluation, drift detection, continual learning.

Self-Correction: Declarative ML is great for baselining and for users who aren’t ML experts, but for complex, high-stakes production systems at FAANG, engineers often need finer-grained control.

NoSQL (Not Only SQL) (Pages 63-65)

  • Movement against relational model’s restrictions (strict schema, schema management pain - #1 reason for Couchbase adoption, footnote 16). SQL can be hard for specialized apps.
  • Many NoSQL systems now also support relational models/SQL.
  • Two major types discussed: Document and Graph.

Document Model (Pages 63-64)

  • Built around “documents” (often JSON, XML, or binary like BSON).
  • Each document has a unique key. Collection of documents ~ table, document ~ row.
  • Flexibility: Documents in a collection can have different schemas (unlike rows in a relational table).
  • “Schemaless” is Misleading: The reading application usually assumes some structure. Responsibility shifts from write-time schema enforcement (relational) to read-time schema interpretation (document).
  • Example (Examples 3-1, 3-2, 3-3): Book data as JSON documents. All info for one book (Harry Potter) is in one document, including “Sold as” array for formats/prices.
  • Better Locality: All info for a book is in one place, easier retrieval than joining multiple relational tables.
  • Worse for Joins/Cross-Document Queries: Finding all books under $25 requires reading all documents, extracting prices, comparing. Less efficient than SQL WHERE price < 25.
  • Many DBs (PostgreSQL, MySQL) now support both relational and document models.

FAANG Perspective: Document databases (MongoDB, DynamoDB) are great for use cases with self-contained data items, flexible schemas, and high scalability needs (e.g., user profiles, product catalogs where attributes vary widely).

Graph Model (Page 65)

  • Data as a graph: nodes and edges (relationships).
  • Prioritizes relationships between data items.
  • Example (Figure 3-5): Social network. Nodes: person, city, country. Edges: lives_in, born_in, coworker, friend, within.
  • Efficient for Relationship-Based Queries: “Find everyone born in USA.” Start at “USA” node, traverse within and born_in edges to find “person” nodes.
  • Hard to do this easily in SQL or document model if #hops is unknown/variable (e.g., 3 hops from Zhenzhong Xu to USA, 2 from Chloe He).

FAANG Perspective: Graph databases (Neo4j, Amazon Neptune) shine for use cases like social networks, knowledge graphs, fraud detection (rings of fraudsters), recommendation (users-who-bought-this-also-bought). Query languages like Cypher or Gremlin are used.

Picking the Right Model: Crucial for simplifying development. Many queries easy in one model are hard in another.

Structured Versus Unstructured Data (Page 66, Table 3-5)

  • Structured Data: Follows a predefined data model/schema (e.g., name=string(50), age=int(0-200)). Easy to analyze (e.g., average age).
    • Disadvantage: Schema changes require retrospective updates, can cause bugs (e.g., new ’email’ field; or null ages becoming 0, confusing ML model - footnote 18’s anecdote, solved by using -1).
  • Unstructured Data: No predefined schema. Usually text, but can be numbers, dates, images, audio (e.g., log files).
    • Advantage: Appealing when business reqs change, or data from many sources can’t conform to one schema.
    • May still have intrinsic patterns (e.g., CSV-like log lines: Lisa,43). But no guarantee all lines follow it.
  • Storage Options:
    • Schema-enforced storage can only store conforming data.
    • Schema-less storage can store any data (e.g., convert all to bytestrings).
  • Data Warehouse: Repository for structured data (processed, ready to use).
  • Data Lake: Repository for unstructured (or raw) data, often before processing.

Table 3-5 Differences:

FeatureStructuredUnstructured
SchemaClearly definedDoesn’t have to follow a schema
Search/AnalyzeEasy(Implied harder until structure is imposed) Fast arrival
Data HandlingSpecific schema onlyAny source
Schema ChangesLots of troubleWorry shifted to downstream apps
Stored InData warehousesData lakes

FAANG Perspective: The distinction is fluid. “Schema-on-read” (data lakes) vs. “schema-on-write” (warehouses). The trend is towards data lakehouses (e.g., Databricks, Snowflake) combining flexibility of lakes with management features of warehouses. Raw data lands in lake, then curated/structured versions are created.

Pages 67-71: Data Storage Engines and Processing

Data formats/models = interface. Storage engines (databases) = implementation on machines. Two main workload types:

Transactional and Analytical Processing (Pages 67-69)

  • Transaction: Digital world: any action (tweet, ride order, model upload, YouTube watch). Inserted as generated, occasionally updated/deleted.
  • Online Transaction Processing (OLTP):
    • Needs to be fast (low latency) for users. High availability. If system can’t process, transaction fails.
    • Transactional Databases: Designed for OLTP. Often associated with ACID properties (Atomicity, Consistency, Isolation, Durability - definitions on page 68 are standard).
      • Atomicity: All steps succeed or all fail (e.g., payment fails, driver not assigned).
      • Consistency: Transactions follow predefined rules (e.g., valid user).
      • Isolation: Concurrent transactions appear isolated (e.g., two users don’t book same driver simultaneously).
      • Durability: Committed transaction persists despite system failure (e.g., phone dies, ride still coming).
    • Not all need ACID. Some find it too restrictive. BASE (Basically Available, Soft state, Eventual consistency) is an alternative, “even more vague” (Kleppmann, footnote 20).
    • Often row-major (transactions processed as units).
  • Online Analytical Processing (OLAP):
    • For analytical questions (e.g., “average ride price in SF in Sept?”). Requires aggregating columns across many rows.
    • Analytical Databases: Designed for this. Efficient with queries from different viewpoints. Often columnar.
  • OLTP/OLAP are Outdated Terms? (Figure 3-6, Google Trends):
    • Separation was due to tech limits (hard to do both well). This is closing.
    • Transactional DBs handling analytical queries (e.g., CockroachDB).
    • Analytical DBs handling transactional queries (e.g., Apache Iceberg, DuckDB).
  • Traditional OLTP/OLAP: Storage and processing tightly coupled. Often meant same data stored multiple times for different query types.
  • Modern Paradigm: Decouple Storage from Processing (Compute). Data in one place, different processing layers on top. (Google BigQuery, Snowflake, IBM, Teradata - footnote 21).
  • “Online” is Overloaded: Used to mean “internet-connected,” then “in production.” Data world: speed of processing/availability (online, nearline, offline - footnote 22).

FAANG Perspective: The decoupling of storage (e.g., S3, GCS) and compute (e.g., Spark, Presto, BigQuery) is a dominant architecture. It provides flexibility, scalability, and cost-efficiency. ACID is critical for financial transactions, eventual consistency is often fine for social media feeds.

ETL: Extract, Transform, and Load (Pages 70-71, Figure 3-7)

  • Early days: relational data, mostly structured. ETL was data warehousing process. Still relevant for ML.
  • General purpose processing/aggregating data into desired shape/format.
  • Extract: Get data from sources. Validate, reject corrupted/malformatted data. Notify sources of rejected data. Crucial first step.
  • Transform: Meaty part. Join, clean, standardize values (Male/Female vs M/F vs 1/2), transpose, deduplicate, sort, aggregate, derive new features, more validation.
  • Load: Decide how/how often to load transformed data into target (file, DB, warehouse).
  • Figure 3-7: Shows sources (DB, App, Flat files) -> ETL -> Targets (Data warehouse, Feature store, DB).
  • Rise of ELT (Extract, Load, Transform):
    • Internet/hardware boom -> easy to collect massive, evolving data. Schemas changed.
    • Idea: Store all raw data in a data lake first (fast arrival, little pre-processing). Applications pull and process as needed.
    • Problem with ELT as data grows: Inefficient to search massive raw data. (Footnote 23: storage cost is rarely a problem now, but processing cost/time is).
  • Trend: Cloud/standardized infra -> committing to predefined schema becomes feasible again.
  • Hybrid: Data Lakehouse (Databricks, Snowflake). Flexibility of lakes + management of warehouses.

FAANG Perspective: ETL/ELT pipelines are the backbone of our data infrastructure. Building robust, scalable, and maintainable ETLs (often using Spark, Beam, Airflow) is a core data engineering function. Feature stores are becoming common “targets” for ML features.

Pages 72-77: Modes of Dataflow – How Data Moves

In production, data isn’t in one process; it flows between many. How does it pass if processes don’t share memory?

Data Passing Through Databases (Page 72)

  • Easiest way: Process A writes to DB, Process B reads from DB.
  • Limitations:
    • Both processes need access to same DB (infeasible if different companies).
    • DB read/writes can be slow, unsuitable for low-latency apps (most consumer-facing ones).

FAANG Perspective: Used for asynchronous tasks or when latency isn’t paramount. E.g., a batch job updates a model quality table, a dashboard service reads it.

Data Passing Through Services (Request-Driven) (Pages 73-74)

  • Direct network communication. Process A requests data from Process B; B returns it.
  • Service-Oriented Architecture (SOA) / Microservices:
    • Process B is a “service” A can call. B can also call A if A is a service.
    • Can be different companies (e.g., investment firm service calls stock exchange service for prices).
    • Can be components of one app (microservices). Allows independent development, testing, maintenance.
  • ML Example: Ride-Sharing Price Optimization (Lyft):
    • Services: Driver Management (available drivers), Ride Management (requested rides), Price Optimization.
    • Price Optimization service needs data from other two to predict optimal price (supply/demand). It requests this data. (Footnote 24: in practice, might use cached data, refresh periodically).
  • Popular Styles: REST (Representational State Transfer) vs. RPC (Remote Procedure Call):
    • REST: For network requests. Often public APIs.
    • RPC: Make remote call look like local function call. Often internal services in same org/datacenter. (Kleppmann, footnote 25).
    • RESTful = implements REST architecture. HTTP is an implementation, not same as REST (footnote 26).

FAANG Perspective: Microservices are ubiquitous. REST for external/public APIs, gRPC (an RPC framework) very common for internal service-to-service communication due to efficiency and strong typing.

Data Passing Through Real-Time Transport (Event-Driven) (Pages 74-77)

  • Motivation (Ride-sharing example, Figure 3-8): If Price Optimization, Driver Mgmt, Ride Mgmt all need data from each other via requests, it becomes a complex web. With hundreds/thousands of services, this is a bottleneck.
  • Request-driven is synchronous: Target service must be listening. If Driver Mgmt is down, Price Opt. keeps resending, times out. Response lost if Price Opt. goes down. Cascading failures.
  • Broker/Event Bus (Figure 3-9): Services communicate via a central broker.
    • Driver Mgmt makes a prediction, broadcasts it (an event) to broker.
    • Other services wanting this data get it from broker.
  • Technically, a DB can be a broker. But slow for low-latency. So, use in-memory storage for brokering (real-time transports).
  • Event-driven architecture: Better for data-heavy systems. Request-driven for logic-heavy.
  • Common Types:
    • Pub/Sub (Publish-Subscribe): Apache Kafka, Amazon Kinesis.
      • Services publish events to topics.
      • Services subscribe to topics to read events. Producers don’t care about consumers.
      • Retention Policy (Figure 3-10): Events kept in in-memory transport for a period (e.g., 7 days), then deleted or moved to permanent storage (e.g., S3). This is key for Kafka’s design – it’s a durable commit log.
    • Message Queue: Apache RocketMQ, RabbitMQ.
      • Event (message) often has intended consumers. Queue gets message to right consumers.
  • Both Kafka/RabbitMQ are very popular (Figure 3-11, Stackshare). (Footnote 27: Mitch Seymour’s Kafka/otters animation is great!)

FAANG Perspective: Kafka is a cornerstone of many real-time data pipelines for logging, metrics, event sourcing, stream processing. It enables decoupling of services and resilience.

Pages 78-79: Batch Processing Versus Stream Processing

Two paradigms for processing data based on its nature (historical vs. in-flight).

Batch Processing

  • Data in storage (DBs, lakes, warehouses) = historical data.
  • Processed in batch jobs (kicked off periodically, e.g., daily job for average surge charge).
  • Distributed systems like MapReduce, Spark process batch data efficiently.
  • Use in ML: Compute features that change less often (static features). E.g., driver’s overall rating (if hundreds of rides, one more doesn’t change it much day-to-day).

Stream Processing

  • Data in real-time transports (Kafka, Kinesis) = streaming data.
  • Computation on this data. Can be periodic (shorter periods, e.g., every 5 mins) or triggered (e.g., user requests ride -> process stream for available drivers).
  • Low latency: Process data as generated, without writing to DB first.
  • Efficiency:
    • Myth: Less efficient than batch (can’t use Spark/MapReduce). Not always true.
    • Streaming tech (Apache Flink) is scalable, distributed (parallel computation).
  • Strength: Stateful computation. Example: 30-day user engagement trial. Batch: recompute over last 30 days daily. Stream: compute on new day’s data, join with older computation (state). Avoids redundancy.
  • Use in ML: Compute features that change quickly (dynamic features / streaming features). E.g., drivers available right now, rides requested last minute, median price of last 10 rides in area. Essential for optimal real-time predictions.
  • Both Batch and Stream Features Needed: Many problems need both. Need infra to process both and join them for ML models (preview of Chapter 7).
  • Stream Computation Engines:
    • Kafka’s built-in stream processing is limited (various data sources).
    • ML streaming features often need complex queries (joins, aggregations).
    • Need efficient stream processing engines: Apache Flink, KSQL (Kafka SQL), Spark Streaming.
    • Flink, KSQL more recognized, nice SQL abstraction for data scientists.
  • Stream processing is harder: Unbounded data, variable rates/speeds.
  • Argument (Flink maintainers, footnote 28): Batch is a special case of streaming. (i.e., a bounded stream). Stream engines can unify both.

FAANG Perspective: This unification is a powerful trend (e.g., Apache Beam model). Lambda architectures (separate batch/stream paths) are complex; Kappa architectures (all stream) are simpler if feasible. Choosing the right features (static, dynamic, or both) is key for model performance and system complexity.

Pages 79-80: Summary of Chapter 3

This chapter built on Chapter 2’s emphasis on data’s importance. Key takeaways:

  • Data Formats: Choose wisely for future use. Row-major vs. column-major, text vs. binary have pros/cons.
  • Data Models: Relational (SQL), Document, Graph. All widely used, each suited for different tasks. Structured (writer assumes schema) vs. Unstructured (reader assumes schema) is fluid.
  • Storage & Processing: Traditionally coupled (OLTP DBs for transactional, OLAP for analytical). Decoupling storage/compute is the trend. Hybrid DBs emerging.
  • Modes of Dataflow: Databases (slow, simple), Services (request-driven, microservices), Real-time Transports (event-driven, Kafka/RabbitMQ for async, low latency).
  • Batch vs. Stream Processing: Historical data -> batch jobs -> static features. Streaming data -> stream engines -> dynamic features. Stream engines can potentially unify both.

With data systems figured out, next chapter is about collecting data and creating training data!

Interview Questions & Page References (Chapter 3)

As promised, here’s a list of potential interview questions related to Chapter 3, with page numbers for where you can find relevant concepts in the book:

General Data Understanding:

  • “Describe the different types of data sources you might encounter in an ML project and their characteristics.” (p. 50-52)
  • “What are the challenges associated with user-input data? How would you handle them?” (p. 50)
  • “Why is system-generated log data important? What are the challenges in managing it?” (p. 50-51)
  • “What are the privacy considerations for user behavior data?” (p. 51, esp. footnote 3)
  • “Explain the difference between first-party, second-party, and third-party data. What are the trends affecting third-party data?” (p. 52)

Data Formats:

  • “What is data serialization? Name some common data formats and their use cases.” (p. 53, Table 3-1)
  • “Compare and contrast row-major (e.g., CSV) and column-major (e.g., Parquet) data formats. When would you choose one over the other?” (p. 54-55, Figure 3-1)
  • “Why might iterating over pandas DataFrame rows be slow? How does its internal storage format relate to this?” (p. 56, Figure 3-2)
  • “Discuss the trade-offs between text formats (like JSON/CSV) and binary formats (like Parquet).” (p. 57, Figure 3-3)
  • “How would you choose a data format for storing a large dataset intended for analytical queries?” (Implied: Parquet, p. 55, 57)

Data Models:

  • “What is the relational data model? Explain the concept of normalization and its pros/cons.” (p. 59-60, Tables 3-2 to 3-4)
  • “What does it mean for SQL to be a declarative language?” (p. 61)
  • “What is NoSQL? Describe the document model and its advantages/disadvantages compared to the relational model.” (p. 63-64, Examples 3-1 to 3-3)
  • “When would a graph data model be appropriate? Give an example.” (p. 65, Figure 3-5)
  • “Explain the difference between structured and unstructured data. What are data lakes and data warehouses?” (p. 66, Table 3-5)

Data Storage Engines & Processing:

  • “What are OLTP and OLAP? How do their requirements differ?” (p. 67-69)
  • “What are ACID properties? Why are they important for transactional databases?” (p. 68)
  • “Why are the terms OLTP/OLAP becoming outdated? What is the significance of decoupling storage and compute?” (p. 69, Figure 3-6)
  • “Describe the ETL process. What are the key steps?” (p. 70-71, Figure 3-7)
  • “What is ELT, and how does it relate to data lakes?” (p. 71)

Modes of Dataflow:

  • “Describe different ways data can be passed between processes in a production system.” (p. 72)
  • “When is passing data through databases suitable/unsuitable?” (p. 72)
  • “Explain request-driven data passing and its connection to microservices. What are REST and RPC?” (p. 73-74)
  • “What is event-driven architecture for data passing? Describe real-time transports like pub/sub (Kafka) and message queues.” (p. 74-77, Figures 3-8 to 3-11)
  • “Why might an event-driven architecture be preferred over a request-driven one for a system with many services?” (p. 75)

Batch vs. Stream Processing:

  • “Compare batch processing and stream processing. What types of data and features are typically associated with each?” (p. 78)
  • “What are the advantages of stream processing, especially concerning latency and stateful computation?” (p. 78)
  • “Name some stream processing engines. Why are they necessary for complex streaming ML features?” (p. 79)
  • “How can batch and stream processing be combined in an ML system?” (p. 79, though more in Ch7)