Skip to content
phattv.dev
EmailLinkedIn

Fundamentals of Data Engineering: Plan and Build Robust Data Systems - Joe Reis & Matt Housley

— books, self-help — 47 min read

About the book

cover

Personal Summary

  1. Context: I started reading this book when i have zero experience as a data engineer. This book focuses on high level of abstraction, not the daily operations with real world examples or applications and lacks the mention of cloud-based solutions as they have mentioned in the Preface.
  2. The most useful take-away is the chart in chapter 1 about The Data Engineering Lifecycle

Part I. Foundation and Building Blocks

Chapter 1: Data Engineering Described

The Data Engineering Lifecycle

data engineering lifecycle

Evolution of the Data Engineer

The early days: 1980 to 2000, from data warehousing to the web

The dot-com boom spawned a ton of activity in web applications and the backend systems to support them—servers, databases, and storage. Much of the infrastructure was expensive, monolithic, and heavily licensed. The vendors selling these backend systems likely didn’t foresee the sheer scale of the data that web applications would produce.

The early 2000s: The birth of contemporary data engineering

Coinciding with the explosion of data, commodity hardware—such as servers, RAM, disks, and flash drives—also became cheap and ubiquitous. Several innovations allowed distributed computation and storage on massive computing clusters at a vast scale. These innovations started decentralizing and breaking apart traditionally monolithic services. The “big data” era had begun.

The Oxford English Dictionary defines big data as “extremely large data sets that may be analyzed computationally to reveal patterns, trends, and associations, especially relating to human behavior and interactions.” Another famous and succinct description of big data is the three Vs of data: velocity, variety, and volume.

As AWS became a highly profitable growth engine for Amazon, other public clouds would soon follow, such as Google Cloud, Microsoft Azure, and DigitalOcean. The public cloud is arguably one of the most significant innovations of the 21st century and spawned a revolution in the way software and data applications are developed and deployed.

The 2000s and 2010s: Big data engineering

Big data quickly became a victim of its own success. As a buzzword, big data gained popularity during the early 2000s through the mid-2010s. Big data captured the imagination of companies trying to make sense of the ever-growing volumes of data and the endless barrage of shameless marketing from companies selling big data tools and services. Because of the immense hype, it was common to see companies using big data tools for small data problems, sometimes standing up a Hadoop cluster to process just a few gigabytes. It seemed like everyone wanted in on the big data action.

Open source developers, clouds, and third parties started looking for ways to abstract, simplify, and make big data available without the high administrative overhead and cost of managing their clusters, and installing, configuring, and upgrading their open source code. The term big data is essentially a relic to describe a particular time and approach to handling large amounts of data.

The 2020s: Engineering for the data lifecycle

What’s old is new again. While “enterprisey” stuff like data management (including data quality and governance) was common for large enterprises in the pre-big-data era, it wasn’t widely adopted in smaller companies. Now that many of the challenging problems of yesterday’s data systems are solved, neatly productized, and packaged, technologists and entrepreneurs have shifted focus back to the “enterprisey” stuff, but with an emphasis on decentralization and agility, which contrasts with the traditional enterprise command-and-control approach.

Data Engineering and Data Science

Data scientists aren’t typically trained to engineer production-grade data systems, and they end up doing this work haphazardly because they lack the support and resources of a data engineer. In an ideal world, data scientists should spend more than 90% of their time focused on the top layers of the pyramid: analytics, experimentation, and ML. When data engineers focus on these bottom parts of the hierarchy, they build a solid foundation for data scientists to succeed.

Data Engineering Skills and Activities

Nowadays, the data-tooling landscape is dramatically less complicated to manage and deploy. Modern data tools considerably abstract and simplify workflows. As a result, data engineers are now focused on balancing the simplest and most cost-effective, best-of-breed services that deliver value to the business. The data engineer is also expected to create agile data architectures that evolve as new trends emerge.

What are some things a data engineer does not do? A data engineer typically does not directly build ML models, create reports or dashboards, perform data analysis, build key performance indicators (KPIs), or develop software applications. A data engineer should have a good functioning understanding of these areas to serve stakeholders best.

Data Maturity and the Data Engineer

Stage 1: Starting with data

A company getting started with data is, by definition, in the very early stages of its data maturity. The company may have fuzzy, loosely defined goals or no goals. Data architecture and infrastructure are in the very early stages of planning and development. Adoption and utilization are likely low or nonexistent. The data team is small, often with a headcount in the single digits. At this stage, a data engineer is usually a generalist and will typically play several other roles, such as data scientist or software engineer. A data engineer’s goal is to move fast, get traction, and add value.

The practicalities of getting value from data are typically poorly understood, but the desire exists. Reports or analyses lack formal structure, and most requests for data are ad hoc. While it’s tempting to jump headfirst into ML at this stage, we don’t recommend it. We’ve seen countless data teams get stuck and fall short when they try to jump to ML without building a solid data foundation.

A data engineer should focus on the following in organizations getting started with data:

  • Get buy-in from key stakeholders, including executive management.
  • Define the right data architecture.
  • Identify and audit data that will support key initiatives.
  • Build a solid data foundation for future data analysts and data scientists to generate reports and models that provide competitive value.

Stage 2: Scaling with data

At this point, a company has moved away from ad hoc data requests and has formal data practices. Now the challenge is creating scalable data architectures and planning for a future where the company is genuinely data-driven. Data engineering roles move from generalists to specialists, with people focusing on particular aspects of the data engineering lifecycle.

In organizations that are in stage 2 of data maturity, a data engineer’s goals are to do the following:

  • Establish formal data practices
  • Create scalable and robust data architectures
  • Adopt DevOps and DataOps practices
  • Build systems that support ML
  • Continue to avoid undifferentiated heavy lifting and customize only when a competitive advantage results

Stage 3: Leading with data

At this stage, the company is data-driven. The automated pipelines and systems created by data engineers allow people within the company to do self-service analytics and ML. Introducing new data sources is seamless, and tangible value is derived. Data engineers implement proper controls and practices to ensure that data is always available to the people and systems. Data engineering roles continue to specialize more deeply than in stage 2.

In organizations in stage 3 of data maturity, a data engineer will continue building on prior stages, plus they will do the following:

  • Create automation for the seamless introduction and usage of new data
  • Focus on building custom tools and systems that leverage data as a competitive advantage
  • Focus on the “enterprisey” aspects of data, such as data management (including data governance and quality) and DataOps
  • Deploy tools that expose and disseminate data throughout the organization, including data catalogs, data lineage tools, and metadata management systems
  • Collaborate efficiently with software engineers, ML engineers, analysts, and others
  • Create a community and environment where people can collaborate and speak openly, no matter their role or position

Business Responsibilities

  • Know how to communicate with nontechnical and technical people.
  • Understand how to scope and gather business and product requirements.
  • Understand the cultural foundations of Agile, DevOps, and DataOps.
  • Control costs.
  • Learn continuously.

A successful data engineer always zooms out to understand the big picture and how to achieve outsized value for the business. Communication is vital, both for technical and nontechnical people. We often see data teams succeed based on their communication with other stakeholders; success or failure is rarely a technology issue. Knowing how to navigate an organization, scope and gather requirements, control costs, and continuously learn will set you apart from the data engineers who rely solely on their technical abilities to carry their career.

Technical Responsibilities

Data engineers remain software engineers, in addition to their many other roles.

What languages should a data engineer know?

  • SQL: The most common interface for databases and data lakes.
  • Python: The bridge language between data engineering and data science.
  • JVM languages such as Java and Scala: Prevalent for Apache open source projects such as Spark, Hive, and Druid.
  • bash: The command-line interface for Linux operating systems.

SQL is a powerful tool that can quickly solve complex analytics and data transformation problems. Data engineers also do well to develop expertise in composing SQL with other operations. Data engineers should also learn modern SQL semantics for dealing with JavaScript Object Notation (JSON) parsing and nested data and consider leveraging a SQL management framework such as dbt (Data Build Tool).

Data Engineers Inside an Organization

Data Engineers and Other Technical Roles

The data engineer is a hub between:

  • data producers: software engineers, data architects, and DevOps or site-reliability engineers (SREs).
  • data consumers: data analysts, data scientists, and ML engineers.

Chapter 2: The Data Engineering Lifecycle

What Is the Data Engineering Lifecycle?

  • Generation: Source Systems
  • Storage
  • Ingestion
  • Transformation
  • Serving Data

Major Undercurrents Across the Data Engineering Lifecycle

  • Security
  • Data Management
  • DataOps
  • Data Architecture
  • Orchestration
  • Software Engineering

Chapter 3: Designing Good Data Architecture

Principles of Good Data Architecture

  • Choose common components wisely.
  • Plan for failure.
  • Architect for scalability.
  • Architecture is leadership.
  • Always be architecting.
  • Build loosely coupled systems.
  • Make reversible decisions.
  • Prioritize security.
  • Embrace FinOps.

Major Architecture Concepts

  • Domains and Services
  • Distributed Systems, Scalability, and Designing for Failure
  • Tight Versus Loose Coupling: Tiers, Monoliths, and Microservices
  • User Access: Single Versus Multitenant
  • Event-Driven Architecture
  • Brownfield Versus Greenfield Projects

Examples and Types of Data Architecture

  • Data Warehouse
  • Data Lake
  • Convergence, Next-Generation Data Lakes, and the Data Platform
  • Modern Data Stack
  • Lambda Architecture
  • Kappa Architecture
  • The Dataflow Model and Unified Batch and Streaming
  • Architecture for IoT
  • Data Mesh

Chapter 4: Choosing Technologies Across the Data Engineering Lifecycle

The following are some considerations for choosing data technologies across the data engineering lifecycle:

  • Team size and capabilities
  • Speed to market
  • Interoperability
  • Cost optimization and business value
  • Today versus the future: immutable versus transitory technologies
  • Location (cloud, on prem, hybrid cloud, multicloud)
  • Build versus buy
  • Monolith versus modular
  • Serverless versus servers
  • Optimization, performance, and the benchmark wars
  • The undercurrents of the data engineering lifecycle

Part II. The Data Engineering Lifecycle in Depth

Chapter 5: Data Generation in Source Systems

Source Systems: Main Ideas

  • Files and Unstructured Data
  • APIs
  • Application Databases (OLTP Systems)
  • Online Analytical Processing System (OLAP Systems)
  • Change Data Capture (CDC)
  • Logs
  • Database Logs
  • CRUD
  • Insert-Only
  • Messages and Streams
  • Types of Time

Source System Practical Details

Databases

  • Database management system
  • Lookups
  • Query optimizer
  • Scaling and distribution
  • Modeling patterns
  • CRUD
  • Consistency
Relational databases

A relational database management system (RDBMS) is one of the most common application backends. Relational databases were developed at IBM in the 1970s and popularized by Oracle in the 1980s.

Data is stored in a table of relations (rows), and each relation contains multiple fields (columns). Tables are typically indexed by a primary key, a unique field for each row in the table.

Tables can also have various foreign keys—fields with values connected with the values of primary keys in other tables, facilitating joins, and allowing for complex schemas that spread data across multiple tables. In particular, it is possible to design a normalized schema. Normalization is a strategy for ensuring that data in records is not duplicated in multiple places, thus avoiding the need to update states in multiple locations at once and preventing inconsistencies

RDBMS systems are typically ACID compliant. Combining a normalized schema, ACID compliance, and support for high transaction rates makes relational database systems ideal for storing rapidly changing application states.

Nonrelational databases: NoSQL

NoSQL, which stands for not only SQL, refers to a whole class of databases that abandon the relational paradigm. There are numerous flavors of NoSQL database designed for almost any imaginable use case.

A key-value database is a nonrelational database that retrieves records using a key that uniquely identifies each record. This is similar to hash map or dictionary data structures presented in many programming languages but potentially more scalable.

Different types of key-value databases offer a variety of performance characteristics to serve various application needs. For example, in-memory key-value databases are popular for caching session data for web and mobile applications, where ultra-fast lookup and high concurrency are required.

Document stores. In this context, a document is a nested object; we can usually think of each document as a JSON object for practical purposes. Documents are stored in collections and retrieved by key. A collection is roughly equivalent to a table in a relational database.

Document databases generally embrace all the flexibility of JSON and don’t enforce schema or types; this is a blessing and a curse. On the one hand, this allows the schema to be highly flexible and expressive. The schema can also evolve as an application grows. On the flip side, we’ve seen document databases become absolute nightmares to manage and query.

Wide-column. A wide-column database is optimized for storing massive amounts of data with high transaction rates and extremely low latency. These databases can scale to extremely high write rates and vast amounts of data. Specifically, wide-column databases can support petabytes of data, millions of requests per second, and sub-10ms latency. These characteristics have made wide-column databases popular in ecommerce, fintech, ad tech, IoT, and real-time personalization applications.

These databases support rapid scans of massive amounts of data, but they do not support complex queries. They have only a single index (the row key) for lookups. Data engineers must generally extract data and send it to a secondary analytics system to run complex queries to deal with these limitations. This can be accomplished by running large scans for the extraction or employing CDC to capture an event stream.

Graph databases. Graph databases explicitly store data with a mathematical graph structure (as a set of nodes and edges).3 Neo4j has proven extremely popular, while Amazon, Oracle, and other vendors offer their graph database products. Roughly speaking, graph databases are a good fit when you want to analyze the connectivity between elements.

In the parlance of graphs, we store nodes (users in the preceding example) and edges (connections between users). Graph databases support rich data models for both nodes and edges. Depending on the underlying graph database engine, graph databases utilize specialized query languages such as SPARQL, Resource Description Framework (RDF), Graph Query Language (GQL), and Cypher.

Search. A search database is a nonrelational database used to search your data’s complex and straightforward semantic and structural characteristics. Two prominent use cases exist for a search database: text search and log analysis.

Text search involves searching a body of text for keywords or phrases, matching on exact, fuzzy, or semantically similar matches. Log analysis is typically used for anomaly detection, real-time monitoring, security analytics, and operational analytics. Queries can be optimized and sped up with the use of indexes.

Time series. A time series is a series of values organized by time. For example, stock prices might move as trades are executed throughout the day, or a weather sensor will take atmospheric temperatures every minute. Any events that are recorded over time—either regularly or sporadically—are time-series data. A time-series database is optimized for retrieving and statistical processing of time-series data.

Measurement data is generated regularly, such as temperature or air-quality sensors. Event-based data is irregular and created every time an event occurs—for instance, when a motion sensor detects movement.

APIs

REST: At virtually any large company, data engineers will need to deal with the problem of writing and maintaining custom code to pull data from APIs, which requires understanding the structure of the data as provided, developing appropriate data-extraction code, and determining a suitable data synchronization strategy.

GraphQL was created at Facebook as a query language for application data and an alternative to generic REST APIs. Whereas REST APIs generally restrict your queries to a specific data model, GraphQL opens up the possibility of retrieving multiple data models in a single request. This allows for more flexible and expressive queries than with REST. GraphQL is built around JSON and returns data in a shape resembling the JSON query.

Webhooks are a simple event-based data-transmission pattern. The data source can be an application backend, a web page, or a mobile app. When specified events happen in the source system, this triggers a call to an HTTP endpoint hosted by the data consumer.

A remote procedure call (RPC) is commonly used in distributed computing. It allows you to run a procedure on a remote system. gRPC is built around the Protocol Buffers open data serialization standard, also developed by Google. gRPC emphasizes the efficient bidirectional exchange of data over HTTP/2

Data Sharing

Data sharing also streamlines the notion of the data marketplace, available on several popular clouds and data platforms. Data marketplaces provide a centralized location for data commerce, where data providers can advertise their offerings and sell them without worrying about the details of managing network access to data systems.

Data sharing allows units of an organization to manage their data and selectively share it with other units while still allowing individual units to manage their compute and query costs separately, facilitating data decentralization. This facilitates decentralized data management patterns such as data mesh.

Third-Party Data Sources

Why would companies want to make their data available? Data is sticky, and a flywheel is created by allowing users to integrate and extend their application into a user’s application. Greater user adoption and usage means more data, which means users can integrate more data into their applications and data systems. The side effect is there are now almost infinite sources of third-party data.

Direct third-party data access is commonly done via APIs, through data sharing on a cloud platform, or through data download. APIs often provide deep integration capabilities, allowing customers to pull and push data. For example, many CRMs offer APIs that their users can integrate into their systems and applications.

Message Queues and Event-Streaming Platforms

A message queue is a mechanism to asynchronously send data (usually as small individual messages, in the kilobytes) between discrete systems using a publish and subscribe model. Data is published to a message queue and is delivered to one or more subscribers. The subscriber acknowledges receipt of the message, removing it from the queue.

Message queues are a critical ingredient for decoupled microservices and event-driven architectures. Some things to keep in mind with message queues are frequency of delivery, message ordering, and scalability.

Message ordering and delivery. Message queues often apply a fuzzy notion of order and first in, first out (FIFO). Strict FIFO means that if message A is ingested before message B, message A will always be delivered before message B. In practice, messages might be published and received out of order, especially in highly distributed message systems.

Delivery frequency. Messages can be sent exactly once or at least once. If a message is sent exactly once, then after the subscriber acknowledges the message, the message disappears and won’t be delivered again. Messages sent at least once can be consumed by multiple subscribers or by the same subscriber more than once. This is great when duplications or redundancy don’t matter. Ideally, systems should be idempotent. In an idempotent system, the outcome of processing a message once is identical to the outcome of processing it multiple times.

Scalability. The most popular message queues utilized in event-driven applications are horizontally scalable, running across multiple servers. This allows these queues to scale up and down dynamically, buffer messages when systems fall behind, and durably store messages for resilience against failure.

An event-streaming platform is used to ingest and process data in an ordered log of records. In an event-streaming platform, data is retained for a while, and it is possible to replay messages from a past point in time.

Topics. In an event-streaming platform, a producer streams events to a topic, a collection of related events.

Stream partitions are subdivisions of a stream into multiple streams. A good analogy is a multilane freeway. Having multiple lanes allows for parallelism and higher throughput. Messages are distributed across partitions by partition key. Messages with the same partition key will always end up in the same partition.

Fault tolerance and resilience. Event-streaming platforms are typically distributed systems, with streams stored on various nodes. If a node goes down, another node replaces it, and the stream is still accessible.

Undercurrents and Their Impact on Source Systems

Security

  • Data is secure and encrypted?
  • Over public internet, or are you using a virtual private network (VPN)?
  • Keep passwords, tokens, and credentials to the source system securely locked away
  • Do you trust the source system?

Data Management

  • Data governance
  • Data quality
  • Schema
  • Schema
  • Master data management
  • Privacy and ethics
  • Regulatory

DataOps

  • Operational excellence—DevOps, DataOps, MLOps, XOps
  • Automation
  • Observability
  • Incident response

Data Architecture

  • Reliability
  • Durability
  • Availability
  • People

Orchestration

  • Cadence and frequency
  • Common frameworks

Software Engineering

  • Networking
  • Authentication and authorization
  • Access patterns
  • Orchestration
  • Parallelization
  • Deployment

Conclusions

If there’s a stick, there’s also a carrot. Better collaboration with source system teams can lead to higher-quality data, more successful outcomes, and better data products. Create a bidirectional flow of communications with your counterparts on these teams; set up processes to notify of schema and application changes that affect analytics and ML. Communicate your data needs proactively to assist application teams in the data engineering process.

Chapter 6: Storage

Whether data is needed seconds, minutes, days, months, or years later, it must persist in storage until systems are ready to consume it for further processing and transmission

Raw Ingredients of Data Storage

  • Magnetic Disk Drive
  • Solid-State Drive (SSD)
  • Random Access Memory (RAM)
  • Networking and CPU
  • Serialization
  • Compression
  • Caching
Storage Types

Data Storage Systems

  • Single Machine Versus Distributed Storage
  • Eventual Versus Strong Consistency
  • File Storage: Local disk storage, Network-attached storage (NAS), Cloud filesystem services
  • Block Storage: RAID, Storage area network (SAN), Cloud virtualized block storage, Local instance volumes
  • Object Storage:
    • Key-value store for immutable data objects, straightforward to manage and use.
    • Object stores for data engineering applications: Object stores are now the gold standard of storage for data lakes. Object storage is an ideal repository for unstructured data in any format beyond these structured data applications.
    • Object lookup: The object store uses a top-level logical container (a bucket in S3 and GCS) and references objects by key. S3 bucket names must be unique across all of AWS. Keys are unique within a bucket.
    • Object consistency and versioning: When we rewrite an object under an existing key in an object store, we’re essentially writing a brand-new object, setting references from the existing key to the object, and deleting the old object references
    • Storage classes and tiers: Cloud vendors now offer storage classes that discount data storage pricing in exchange for reduced access or reduced durability.
  • Cache and Memory-Based Storage Systems: Memcached & Redis
  • The Hadoop Distributed File System: Hadoop is similar to object storage but with a key difference: Hadoop combines compute and storage on the same nodes, where object stores typically have limited support for internal processing.
  • Streaming Storage
  • Indexes, Partitioning, and Clustering
    • The evolution from rows to columns: Columnar serialization allows a database to scan only the columns required for a particular query, sometimes dramatically reducing the amount of data read from the disk.
    • From indexes to partitions and clustering: In addition to scanning only data in columns relevant to a query, we can partition a table into multiple subtables by splitting it on a field. Clusters allow finer-grained organization of data within partitions

Data Engineering Storage Abstractions

  • The Data Warehouse: Data warehouses are a standard OLAP data architecture
  • The Data Lake: The data lake was originally conceived as a massive store where data was retained in raw, unprocessed form.
  • The Data Lakehouse: The data lakehouse is an architecture that combines aspects of the data warehouse and the data lake. A lakehouse system is a metadata and file-management layer deployed with data management and transformation tools. Databricks has heavily promoted the lakehouse concept with Delta Lake, an open source storage management system.
  • Data Platforms: ecosystems of interoperable tools with tight integration into the core data storage layer.
  • Stream-to-Batch Storage Architecture: The query engine supports seamless querying of both the streaming buffer and the object data to provide users a current, nearly real-time view of the table.
  • Data Catalog: A data catalog is a centralized metadata store for all data across an organization.
  • Data Sharing: Data sharing allows organizations and individuals to share specific data and carefully defined permissions with specific entities.
  • Schema: data becomes more useful when we have as much information about its structure and organization.
  • Separation of Compute from Storage: AWS EMR with S3 and HDFS, Apache Spark, Apache Druid, Hybrid object storage
  • Data Storage Lifecycle and Data Retention: what data do you need to keep, and how long should you keep it. Data Storage Lifecycle
  • Single-Tenant Versus Multitenant Storage: Multitenant storage allows for the storage of multiple tenants within a single database. For example, instead of the single-tenant scenario where customers get their own database, multiple customers may reside in the same database schemas or tables in a multitenant database

Undercurrents

  • Security
  • Data Management
  • DataOps
  • Data Architecture
  • Orchestration
  • Software Engineering

Chapter 7: Ingestion

What Is Data Ingestion?

Data ingestion is the process of moving data from one place to another. Data ingestion implies data movement from source systems into storage in the data engineering lifecycle, with ingestion as an intermediate step

Key Engineering Considerations for the Ingestion Phase

What’s the use case for the data I’m ingesting?

  • Can I reuse this data and avoid ingesting multiple versions of the same dataset?
  • Where is the data going? What’s the destination?
  • How often should the data be updated from the source?
  • What is the expected data volume?
  • What format is the data in? Can downstream storage and transformation accept this format?
  • Is the source data in good shape for immediate downstream use? That is, is the data of good quality? What post-processing is required to serve it? What are data-quality risks (e.g., could bot traffic to a website contaminate the data)?
  • Does the data require in-flight processing for downstream ingestion if the data is from a streaming source?

Bounded Versus Unbounded Data

Unbounded data is data as it exists in reality, as events happen, either sporadically or continuously, ongoing and flowing. Bounded data is a convenient way of bucketing data across some sort of boundary, such as time.. All data is unbounded until it’s bounded.

Frequency

Ingestion Frequency

Synchronous Versus Asynchronous Ingestion

With synchronous ingestion, the source, ingestion, and destination have complex dependencies and are tightly coupled. With asynchronous ingestion, dependencies can now operate at the level of individual events, much as they would in a software backend built from microservices.

Serialization and Deserialization

Serialization means encoding the data from a source and preparing data structures for transmission and intermediate storage stages.

Throughput and Scalability

In theory, your ingestion should never be a bottleneck. In practice, ingestion bottlenecks are pretty standard. Data throughput and system scalability become critical as your data volumes grow and requirements change. Where you’re ingesting data from matters a lot. Another thing to consider is your ability to handle bursty data ingestion.

Reliability and Durability

Reliability entails high uptime and proper failover for ingestion systems. Durability entails making sure that data isn’t lost or corrupted. Continually evaluate the trade-offs and costs of reliability and durability.

Payload

  • Kind: data has a type—tabular, image, video, text, etc. The type directly influences the data format or the way it is expressed in bytes, names, and file extensions
  • Shape: tabular, semistructured json, unstructured text, images, uncompressed audio
  • Size: the size of the data describes the number of bytes of a payload
  • Schema and data types: many data payloads have a schema, such as tabular and semistructured data. Other data, such as unstructured text, images, and audio, will not have an explicit schema or data types.
  • Detecting and handling upstream and downstream schema changes: engineers must still implement strategies to respond to changes automatically and alert on changes that cannot be accommodated automatically.
  • Schema registries: a schema registry is a metadata repository used to maintain schema and data type integrity in the face of constantly changing schemas
  • Metadata: metadata is data about data. Metadata can be as critical as the data itself.

Push Versus Pull Versus Poll Patterns

  • A push strategy involves a source system sending data to a target, while a pull strategy entails a target reading data directly from a source. Another pattern related to pulling is polling for data.
  • Polling involves periodically checking a data source for any changes. When changes are detected, the destination pulls the data as it would in a regular pull situation.

Batch Ingestion Considerations

Batch ingestion, which involves processing data in bulk, is often a convenient way to ingest data. This means that data is ingested by taking a subset of data from a source system, based either on a time interval or the size of accumulated data.

  • Snapshot or Differential Extraction: Data engineers must choose whether to capture full snapshots of a source system or differential (sometimes called incremental) updates
  • File-Based Export and Ingestion: With file-based ingestion, export processes are run on the data-source side, giving source system engineers complete control over what data gets exported and how the data is preprocessed.
  • ETL Versus ELT
  • Inserts, Updates, and Batch Size: Batch-oriented systems often perform poorly when users attempt to perform many small-batch operations rather than a smaller number of large operations.
  • Data Migration

Message and Stream Ingestion Considerations

  • Schema Evolution
  • Late-Arriving Data
  • Ordering and Multiple Delivery
  • Replay
  • Time to Live
  • Message Size
  • Error Handling and Dead-Letter Queues
  • Consumer Pull and Push
  • Location

Ways to Ingest Data

  • Direct Database Connection
  • Change Data Capture: Batch-oriented, Continuous & database replication
  • APIs
  • Message Queues and Event-Streaming Platforms
  • Managed Data Connectors
  • Moving Data with Object Storage
  • EDI
  • Databases and File Export
  • Practical Issues with Common File Formats
  • Shell
  • SSH
  • SFTP and SCP
  • Webhooks
  • Web Interface
  • Web Scraping
  • Transfer Appliances for Data Migration
  • Data Sharing

Chapter 8: Queries, Modeling, and Transformation

Queries

What Is a Query?

A query allows you to retrieve and act on data.

  • Data definition language: Data engineers use common SQL DDL expressions: CREATE, DROP, and UPDATE.
  • Data manipulation language: SELECT, INSERT, UPDATE, DELETE, COPY, MERGE.
  • Data control language: GRANT, DENY, and REVOKE.
  • Transaction control language: COMMIT and ROLLBACK.

The Life of a Query

SQL Query Lifecycle

The Query Optimizer

A query optimizer’s job is to optimize query performance and minimize costs by breaking the query into appropriate steps in an efficient order. The optimizer will assess joins, indexes, data scan size, and other factors. The query optimizer attempts to execute the query in the least expensive manner.

Improving Query Performance

  • Optimize your join strategy and schema: prejoin, consider the details and complexity of the join condition, use common table expressions (CTE)
  • Use the explain plan and understand your query’s performance: In addition to using EXPLAIN to understand how your query will run, you should monitor your query’s performance, viewing metrics on database resource consumption.
  • Avoid full table scans
  • Know how your database handles commits: You should be intimately familiar with how your database handles commits and transactions, and determine the expected consistency of query results.
  • Vacuum dead records: As these old records accumulate in the database filesystem, they eventually no longer need to be referenced. You should remove these dead records in a process called vacuuming. Vacuuming becomes even more critical for relational databases such as PostgreSQL and MySQL.
    • First, it frees up space for new records, leading to less table bloat and faster queries.
    • Second, new and relevant records mean query plans are more accurate; outdated records can lead the query optimizer to generate suboptimal and inaccurate plans.
    • Finally, vacuuming cleans up poor indexes, allowing for better index performance.
  • Leverage cached query results

Queries on Streaming Data

  • Basic query patterns on streams
  • The fast-follower approach
  • The Kappa architecture: handle all data like events and store these events as a stream rather than a table
  • Windows, triggers, emitted statistics, and late-arriving data
  • Session window
  • Fixed-time windows
  • Sliding windows
  • Watermarks
  • Combining streams with other data
  • Conventional table joins
  • Enrichment

Stream-to-stream joining

Data Modeling

What Is a Data Model?

A data model represents the way data relates to the real world. A good data model captures how communication and work naturally flow within your organization. In contrast, a poor data model (or nonexistent one) is haphazard, confusing, and incoherent.

Conceptual, Logical, and Physical Data Models

  • Conceptual: Contains business logic and rules and describes the system’s data, such as schemas, tables, and fields (names and types).
  • Logical: Details how the conceptual model will be implemented in practice by adding significantly more detail.
  • Physical: Defines how the logical model will be implemented in a database system.

Normalization

Normalization is a database data modeling practice that enforces strict control over the relationships of tables and columns within a database. The goal of normalization is to remove the redundancy of data within a database and ensure referential integrity. Basically, it’s don’t repeat yourself (DRY) applied to data in a database.

  • Denormalized: No normalization. Nested and redundant data is allowed.
  • First normal form (1NF): Each column is unique and has a single value. The table has a unique primary key.
  • Second normal form (2NF): The requirements of 1NF, plus partial dependencies are removed.
  • Third normal form (3NF): The requirements of 2NF, plus each table contains only relevant fields related to its primary key and has no transitive dependencies.

Techniques for Modeling Batch Analytical Data

  • Inmon: The four critical parts of a data warehouse can be described as follows:
    • Subject-oriented: The data warehouse focuses on a specific subject area, such as sales or marketing.
    • Integrated: Data from disparate sources is consolidated and normalized.
    • Nonvolatile: Data remains unchanged after data is stored in a data warehouse.
    • Time-variant: Varying time ranges can be queried.
  • Kimball: data is modeled with two general types of tables: facts and dimensions. You can think of a fact table as a table of numbers, and dimension tables as qualitative data referencing a fact. Dimension tables surround a single fact table in a relationship called a star schema.
  • Fact tables: contain factual, quantitative, and event-related data. The data in a fact table is immutable because facts relate to events. Therefore, fact tables don’t change and are append-only. Fact tables are typically narrow and long, meaning they have not a lot of columns but a lot of rows that represent events. Fact tables should be at the lowest grain possible.
  • Dimensions tables: Dimension tables provide the reference data, attributes, and relational context for the events stored in fact tables. Dimension tables are smaller than fact tables and take an opposite shape, typically wide and short.
  • Star schema: Unlike highly normalized approaches to data modeling, the star schema is a fact table surrounded by the necessary dimensions. This results in fewer joins than other data models, which speeds up query performance. Another advantage of a star schema is it’s arguably easier for business users to understand and use.
  • Data Vault: A Data Vault model consists of three main types of tables: hubs, links, and satellites. In short, a hub stores business keys, a link maintains relationships among business keys, and a satellite represents a business key’s attributes and context.
  • Wide denormalized tables: The wide table simply contains all of the data you would have joined in a more rigorous modeling approach. Facts and dimensions are represented in the same table. The lack of data model rigor also means not a lot of thought is involved. Load your data into a wide table and start querying it.

Modeling Streaming Data

The streaming data experts we’ve talked with overwhelmingly suggest you anticipate changes in the source data and keep a flexible schema. This means there’s no rigid data model in the analytical database. Instead, assume the source systems are providing the correct data with the right business definition and logic, as it exists today. And because storage is cheap, store the recent streaming and saved historical data in a way they can be queried together. Optimize for comprehensive analytics against a dataset with a flexible schema.

Transformations

Transformations manipulate, enhance, and save data for downstream use, increasing its value in a scalable, reliable, and cost-effective manner.

A transformation differs from a query. A query retrieves the data from various sources based on filtering and join logic. A transformation persists the results for consumption by additional transformations or queries. These results may be stored ephemerally or permanently.

Besides persistence, a second aspect that differentiates transformations from queries is complexity. You’ll likely build complex pipelines that combine data from multiple sources and reuse intermediate results for multiple final outputs. These complex pipelines might normalize, model, aggregate, or featurize data.

Batch Transformations

  • Distributed joins: break a logical join (the join defined by the query logic) into much smaller node joins that run on individual servers in the cluster.
  • Broadcast join: A broadcast join is generally asymmetric, with one large table distributed across nodes and one small table that can easily fit on a single node.
  • Shuffle hash join: If neither table is small enough to fit on a single node, the query engine will use a shuffle hash join.
  • ETL, ELT, and data pipelines: Organizations no longer need to standardize on ETL or ELT but can instead focus on applying the proper technique on a case-by-case basis as they build data pipelines.
  • SQL and code-based transformation tools.
  • SQL is declarative...but it can still build complex data workflows.
  • Update patterns
  • Truncate and reload
  • Insert only
  • Delete: A hard delete permanently removes a record from a database, while a soft delete marks the record as “deleted.”
  • Upsert/merge: Upserting takes a set of source records and looks for matches against a target table by using a primary key or another logical condition.
  • Schema updates
  • Data wrangling: Data wrangling takes messy, malformed data and turns it into useful, clean data. This is generally a batch transformation process.
  • Business logic and derived data: One of the most common use cases for transformation is to render business logic.
  • MapReduce: A simple MapReduce job consists of a collection of map tasks that read individual data blocks scattered across the nodes, followed by a shuffle that redistributes result data across the cluster and a reduce step that aggregates data on each node.
  • After MapReduce: The cloud is one of the drivers for the broader adoption of memory caching; it is much more effective to lease memory during a specific processing job than to own it 24 hours a day. Advancements in leveraging memory for transformations will continue to yield gains for the foreseeable future.

Materialized Views, Federation, and Query Virtualization

  • Views: A view is a database object that we can select from just like any other table. In practice, a view is just a query that references other tables. When we select from a view, that database creates a new query that combines the view subquery with our query. The query optimizer then optimizes and runs the full query.
  • Materialized views: A materialized view does some or all of the view computation in advance.
  • Composable materialized views: Databricks has introduced the notion of live tables. Each table is updated as data arrives from sources. Data flows down to subsequent tables asynchronously.
  • Federated queries: Federated queries are a database feature that allows an OLAP database to select from an external data source, such as object storage or RDBMS.
  • Data virtualization: Data virtualization can be viewed as a tool that expands the data lake to many more sources by abstracting away barriers used to silo data between organizational units

Streaming Transformations and Processing

  • Basics: Streaming transformations aim to prepare data for downstream consumption.
  • Transformations and queries are a continuum
  • Streaming DAGs
  • Micro-batch versus true streaming

Chapter 9: Serving Data for Analytics, Machine Learning, and Reverse ETL

General Considerations for Serving Data

  • Trust: Above all else, trust is the root consideration in serving data; end users need to trust the data they’re receiving. To realize data quality and build stakeholder trust, utilize data validation and data observability processes, in conjunction with visually inspecting and confirming validity with stakeholders
  • What’s the Use Case, and Who’s the User? Data is at its best when it leads to action. Always approach data engineering from the perspective of the user and their use case.
  • Data Products: A good data product has positive feedback loops. More usage of a data product generates more useful data, which is used to improve the data product. Rinse and repeat.
  • Self-Service or Not?
  • Data Definitions and Logic: Data definition refers to the meaning of data as it is understood throughout the organization. Data logic stipulates formulas for deriving metrics from data—say, gross sales or customer lifetime value.
  • Data Mesh: Instead of siloed data teams serving their internal constituents, every domain team takes on two aspects of decentralized, peer-to-peer data serving.

Analytics

The first data-serving use case you’ll likely encounter is analytics, which is discovering, exploring, identifying, and making visible key insights and patterns within data.

Business Analytics

Business analytics uses historical and current data to make strategic and actionable decisions. The types of decisions tend to factor in longer-term trends and often involve a mix of statistical and trend analysis, alongside domain expertise and human judgment.

A dashboard concisely shows decision makers how an organization is performing against a handful of core metrics, such as sales and customer retention.

The goal of a report is to use data to drive insights and action. The analyst runs some SQL queries in the data warehouse, aggregates the return codes that customers provide as the reason for their return, and discovers that the fabric in the running shorts is of inferior quality, often wearing out within a few uses

The analyst was asked to dig into a potential issue and come back with insights. This represents an example of ad hoc analysis. Reports typically start as ad hoc requests. If the results of the ad hoc analysis are impactful, they often end up in a report or dashboard.

Operational Analytics

Operational analytics versus business analytics = immediate action versus actionable insights.

Operational analytics is quite the opposite, as real-time updates can be impactful in addressing a problem when it occur.

An example of operational analytics is real-time application monitoring. Many software engineering teams want to know how their application is performing; if issues arise, they want to be notified immediately.

Embedded Analytics

Whereas business and operational analytics are internally focused, a recent trend is external-facing or embedded analytics. With so much data powering applications, companies increasingly provide analytics to end users.

Machine Learning

The second major area for serving data is machine learning.

In some organizations, ML engineers take over data processing for ML applications right after data collection or may even form an entirely separate and parallel data organization that handles the entire lifecycle for all ML applications. Data engineers handle all data processing in other settings and then hand off data to ML engineers for model training. Data engineers may even handle some extremely ML-specific tasks, such as featurization of data.

What a Data Engineer Should Know About ML

  • supervised, unsupervised, and semisupervised learning.
  • classification and regression techniques.
  • various techniques for handling time-series data. This includes time-series analysis, as well as time-series forecasting.
  • When to use the “classical” techniques (logistic regression, tree-based learning, support vector machines) versus deep learning.
  • When would you use automated machine learning (AutoML) versus handcrafting an ML model>
  • What are data-wrangling techniques used for structured and unstructured data?
  • If you’re serving structured or semistructured data, ensure that the data can be properly converted during the feature-engineering process.
  • How to encode categorical data and the embeddings for various types of data.
  • The difference between batch and online learning. Which approach is appropriate for your use case?
  • How does the data engineering lifecycle intersect with the ML lifecycle at your company? Will you be responsible for interfacing with or supporting ML technologies such as feature stores or ML observability?
  • Know when it’s appropriate to train locally, on a cluster, or at the edge. When would you use a GPU over a CPU? The type of hardware you use largely depends on the type of ML problem you’re solving, the technique you’re using, and the size of your dataset.
  • Know the difference between the applications of batch and streaming data in training ML models.
  • What are data cascades, and how might they impact ML models?
  • Are results returned in real time or in batch?
  • The use of structured versus unstructured data. We might cluster tabular (structured) customer data or recognize images (unstructured) by using a neural net.

Ways to Serve Data for Analytics and ML

  • File Exchange: File exchange is ubiquitous in data serving. We process data and generate files to pass to data consumers.
  • Databases: Serving data from a database carries a variety of benefits. A database imposes order and structure on the data through schema.
  • Streaming Systems: At a high level, understand that this type of serving may involve emitted metrics, which are different from traditional queries.
  • Query Federation: Instead of serving data from a single system, you’re now serving data from multiple systems, each with its usage patterns, quirks, and nuances. If federated queries touch live production source systems, you must ensure that the federated query won’t consume excessive resources in the source.
  • Data Sharing
  • Semantic and Metrics Layers
  • Serving Data in Notebooks

Reverse ETL

Reverse ETL

Reverse ETL takes processed data from the output side of the data engineering lifecycle and feeds it back into source systems

Conclusion

The data engineering lifecycle has a logical ending at the serving stage. As with all lifecycles, a feedback loop occurs. You should view the serving stage as a chance to learn what’s working and what can be improved. Listen to your stakeholders. If they bring up issues—and they inevitably will—try not to take offense. Instead, use this as an opportunity to improve what you’ve built.

A good data engineer is always open to new feedback and constantly finds ways to improve their craft.

Part III. Security, Privacy, and the Future of Data Engineering

Chapter 10: Security and Privacy

Security is a key ingredient for privacy. Privacy has long been critical to trust in the corporate information technology space; engineers directly or indirectly handle data related to people’s private lives.

People

The weakest link in security and privacy is you. Security is often compromised at the human level, so conduct yourself as if you’re always a target. A bot or human actor is trying to infiltrate your sensitive credentials and information at any given time.

  • The Power of Negative Thinking: Positive thinking can blind us to the possibility of terrorist attacks or medical emergencies and deter preparation. Negative thinking allows us to consider disastrous scenarios and act to prevent them. The best way to protect private and sensitive data is to avoid ingesting this data in the first place.
  • Always Be ParanoidAlways exercise caution when someone asks you for your credentials. When in doubt—and you should always be in extreme doubt when asked for credentials— hold off and get second opinions from your coworkers and friends.

Processes

When people follow regular security processes, security becomes part of the job. Make security a habit, regularly practice real security, exercise the principle of least privilege, and understand the shared responsibility model in the cloud.

Security Theater Versus Security Habit: Security needs to be simple and effective enough to become habitual throughout an organization.

  • Active SecurityReturning to the idea of negative thinking, active security entails thinking about and researching security threats in a dynamic and changing world.
  • The principle of least privilege: means that a person or system should be given only the privileges and data they need to complete the task at hand and nothing more.
  • Shared Responsibility in the Cloud
  • Always Back Up Your Data: You need to back up your data regularly, both for disaster recovery and continuity of business operations, if a version of your data is compromised in a ransomware attack.

Technology

  • Patch and Update Systems: To avoid exposing a security flaw in an older version of the tools you’re using, always patch and update operating systems and software as new updates become available.
  • Encryption: Encryption is a baseline requirement for any organization that respects security and privacy. It will protect you from basic attacks, such as network traffic interception.
  • Logging, Monitoring, and Alerting: Most companies don’t find out about security incidents until well after the fact. Part of DataOps is to observe, detect, and alert on incidents. As a data engineer, you should set up automated monitoring, logging, and alerting to be aware of peculiar events when they happen in your systems. If possible, set up automatic anomaly detection.
  • Network Access
  • Security for Low-Level Data Engineering

Conclusion

Security needs to be a habit of mind and action; treat data like your wallet or smartphone.

Chapter 11: The Future of Data Engineering

  • The Data Engineering Lifecycle Isn’t Going Away: As companies realize they first need to build a data foundation before moving to “sexier” things like AI and ML, data engineering will continue growing in popularity and importance. This progress centers around the data engineering lifecycle.
  • The Decline of Complexity and the Rise of Easy-to-Use Data Tools: Simplified, easy-to-use tools continue to lower the barrier to entry for data engineering. The trend toward simplicity will continue. Data engineering isn’t dependent on a particular technology or data size. Data engineering is now something that all companies can do.
  • The Cloud-Scale Data OS and Improved Interoperability: We will also see significant improvements in the scaffolding that manages cloud data services. We will also see significant enhancements in the domain of live data.
  • “Enterprisey” Data Engineering: This allows data engineers working on new tooling to find opportunities in the abstractions of data management, DataOps, and all the other undercurrents of data engineering. Data engineers will become “enterprisey.”
  • Titles and Responsibilities Will Morph: While the data engineering lifecycle isn’t going anywhere anytime soon, the boundaries between software engineering, data engineering, data science, and ML engineering are increasingly fuzzy. As simplicity moves up the stack, data scientists will spend a smaller slice of their time gathering and munging data. As data becomes more tightly embedded in every business’s processes, new roles will emerge in the realm of data and algorithms. Another area in which titles may morph is at the intersection of software engineering and data engineering.
  • Moving Beyond the Modern Data Stack, Toward the Live Data Stack.
    • Streaming Pipelines and Real-Time Analytical Databases: Streaming technologies will continue to see extreme growth for the foreseeable future. This will happen in conjunction with a clearer focus on the business utility of streaming data. Real-time analytical databases enable both fast ingestion and subsecond queries on this data. This data can be enriched or combined with historical datasets. When combined with a streaming pipeline and automation, or dashboard that is capable of real-time analytics, a whole new level of possibilities opens up.
    • The Fusion of Data with Applications: Soon, application stacks will be data stacks, and vice versa. Applications will integrate real-time automation and decision making, powered by the streaming pipelines and ML. The data engineering lifecycle won’t necessarily change, but the time between stages of the lifecycle will drastically shorten. A lot of innovation will occur in new technologies and practices that will improve the experience of engineering the live data stack.
    • The Tight Feedback Between Applications and ML: High volumes of fast-moving data, coupled with sophisticated workflows and actions, are candidates for ML. As data feedback loops become shorter, we expect most applications to integrate ML. As data moves more quickly, the feedback loop between applications and ML will tighten.
    • Dark Matter Data and the Rise of...Spreadsheets?! What’s the most widely used data platform? It’s the humble spreadsheet. Depending on the estimates you read, the user base of spreadsheets is between 700 million and 2 billion people. Spreadsheets are the dark matter of the data world. A good deal of data analytics runs in spreadsheets and never makes its way into the sophisticated data systems that we describe in this book. In many organizations, spreadsheets handle financial reporting, supply-chain analytics, and even CRM.

Conclusion

Data engineering is a vast topic; while we could not go into any technical depth in individual areas, we hope that we have succeeded in creating a kind of travel guide that will help current data engineers, future data engineers, and those who work adjacent to the field to find their way in a domain that is in flux. We advise you to continue exploration on your own.