Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.2 The eight considerations when building data pipelines

This is exactly the workaround Ch. 8 §3.6 described ("manage offsets in the database") — Connect gives connectors a first-class hook for it.

"Those challenges are not specific to Kafka but are general data integration problems."

2.1 Timeliness

"Some systems expect their data to arrive in large bulks once a day; others expect the data to arrive a few milliseconds after it is generated. Most data pipelines fit somewhere in between."

What a good system must do: "support different timeliness requirements for different pipelines and also make the MIGRATION between different timetables easier as business requirements change."

producerswrite as frequently or infrequently as neededKafka — a giant bufferdecouples the time-sensitivity requirementsreal-time consumerread events as they arrivebatch consumerrun every hour, read what accumulatedKafka spans the whole range with one substrate — the same topic serves both consumers.

“Producers can write events in real time, while consumers process batches of events, or vice versa.” The batch consumer's contract is literally “run every hour, connect to Kafka, and read the events that accumulated during the previous hour.”

Figure 9.2.12.1 Timeliness

And back pressure comes free:

"This also makes it trivial to apply back pressure — Kafka itself applies back pressure on producers (by delaying acks when needed) since consumption rate is driven ENTIRELY by the consumers."

(Note the mechanism: Kafka's pull model means a slow consumer can never stall a producer directly — Ch. 1's ActiveMQ lesson — but the broker can still throttle producers via delayed acks and quotas when it needs to.)

2.2 Reliability

"We want to avoid single points of failure and allow for fast and automatic recovery from all sorts of failure events. Data pipelines are often the way data arrives to business-critical systems; failure for more than a few seconds can be hugely disruptive, especially when the timeliness requirement is closer to the few-milliseconds end."

Delivery guarantees:

at-least-onceexactly-oncesource systemdestinationduplicates possiblesource systemdestinationno loss, no duplication
  • at-least-once — “every event from the source system will reach its destination, but sometimes retries will cause duplicates”
  • exactly-once — “every event from the source system will reach the destination with no possibility for loss or duplication”
Figure 9.2.2Delivery guarantees

Where Kafka stands (per Ch. 7/8):

"Kafka can provide at-least-once on its own, and exactly-once when combined with an external data store that has a TRANSACTIONAL MODEL or UNIQUE KEYS. Since many of the end points ARE data stores that provide the right semantics for exactly-once delivery, a Kafka-based pipeline can often be implemented as exactly-once."

"It is worth highlighting that Kafka's Connect API makes it easier for connectors to build an end-to-end exactly-once pipeline by providing an API for integrating with the external systems WHEN HANDLING OFFSETS. Indeed, many of the available open source connectors support exactly-once delivery."

This is exactly the workaround Ch. 8 §3.6 described ("manage offsets in the database") — Connect gives connectors a first-class hook for it.

2.3 High and varying throughput

"they should be able to scale to very high throughputs... Even more importantly, they should be able to ADAPT IF THROUGHPUT SUDDENLY INCREASES."

Without a buffer
back-pressure neededproducerconsumer

Consumer throughput must be coupled to producer throughput, so the pipeline needs a complex back-pressure mechanism.

With Kafka as the buffer
producerKafkadata accumulates hereconsumer
  • “We no longer need to couple consumer throughput to producer throughput.”
  • “If producer throughput exceeds that of the consumer, data will accumulate in Kafka until the consumer can catch up.”
  • ⇒ scale either side “dynamically and independently”.
Figure 9.2.32.3 High and varying throughput

Kafka's own numbers: "capable of processing hundreds of megabytes per second on even modest clusters."

Connect's parallelism story: "the Kafka Connect API focuses on parallelizing the work and can do this on a single node as well as by scaling out, depending on system requirements... allows data sources and sinks to split the work among multiple threads of execution and use the available CPU resources EVEN WHEN RUNNING ON A SINGLE MACHINE."

Plus: "Kafka also supports several types of compression, allowing users and admins to control the use of network and storage resources."

2.4 Data formats

"One of the most important considerations... is reconciling different data formats and data types."

The realistic mess the book describes:

XML + relational dataKafkaas AvroElasticsearchas JSONHDFSas ParquetS3as CSV
Figure 9.2.4The realistic mess the book describes

Kafka's stance — deliberate agnosticism:

"Kafka itself and the Connect API are COMPLETELY AGNOSTIC when it comes to data formats. ... Kafka Connect has its own in-memory objects that include data types and schemas, but... it allows for pluggable CONVERTERS to allow storing these records in any format. This means that NO MATTER WHICH DATA FORMAT YOU USE FOR KAFKA, IT DOES NOT RESTRICT YOUR CHOICE OF CONNECTORS."

Schema propagation — the aspirational behavior:

"Many sources and sinks have a schema; we can read the schema from the source with the data, store it, and use it to validate compatibility or even UPDATE THE SCHEMA IN THE SINK DATABASE. A classic example is a data pipeline from MySQL to Snowflake. If someone added a column in MySQL, a great pipeline will make sure the column gets added to Snowflake too as we are loading new data into it."

Sink-side format is the connector's job: "sink connectors are responsible for the format in which the data is written to the external system. Some connectors choose to make this format pluggable — for example, the S3 connector allows a choice between Avro and Parquet."

Behavioral differences, not just format differences:

pushesframework must pullappendappend / updateSyslogrelational databasesdata integration frameworkHDFSappend-only — we can only writemost systemsappend and update

► “A generic data integration framework should also handle differences in behavior between various sources and sinks.”

Figure 9.2.5Behavioral differences, not just format differences

2.5 Transformations — ETL vs ELT

"Transformations are more controversial than other requirements."

Extractsource systemTransformin the pipeline, in flightLoadtarget systemThe pipeline is responsible for making modifications to the data as it passes through.
  • ✓ “perceived benefit of saving time and storage because you don't need to store the data, modify it, and store it again”
  • ⚠ “Depending on the transformations, this benefit is sometimes real, but sometimes it shifts the burden of computation and storage to the data pipeline itself, which may or may not be desirable.”
  • ✗ The main drawback: “the transformations that happen to the data in the pipeline may tie the hands of those who wish to process the data further down the pipe.”
Figure 9.2.62.5 Transformations — ETL vs ELT

The worked cost of ETL's drawback:

"If the person who built the pipeline between MongoDB and MySQL decided to filter certain events or remove fields, all the users and applications who access the data in MySQL will only have access to partial data. If they require access to the missing fields, the pipeline needs to be REBUILT, and historical data will require REPROCESSING (assuming it is available)."

That parenthetical — "assuming it is available" — is the real danger. A dropped field is often unrecoverable.

done at the targetExtractLoadminimal transformationtarget systemcollects “raw data”Transformall required processingMinimal transformation is mostly data type conversion — the arriving data stays as similaras possible to the source data.
  • ✓ “maximum flexibility to users of the target system, since they have access to all the data”
  • ✓ “easier to troubleshoot since all data processing is limited to one system rather than split between the pipeline and additional apps”
  • ✗ “the transformations take CPU and storage resources at the target system. In some cases, these systems are expensive and there is strong motivation to move computation off those systems.”
Figure 9.2.7

Kafka's split:

The transformationWhere it runs
Stateless, in-flightKafka Connect single message transformations (SMTs) — “routing messages to different topics, filtering messages, changing data types, redacting specific fields, and more”
Stateful — joins, aggregationsKafka Streams (Ch. 14)

⚠️ WARNING — the one-to-many rule

*"When building an ETL system with Kafka, keep in mind that Kafka allows you to build ONE-TO-MANY pipelines, where the source data is written to Kafka once and then consumed by multiple applications and written to multiple target systems.

Some preprocessing and cleanup IS expected, such as:

  • standardizing timestamps and data types
  • adding lineage
  • perhaps removing personal information

— transformations that will benefit ALL consumers of the data.

But DON'T PREMATURELY CLEAN AND OPTIMIZE THE DATA ON INGEST BECAUSE IT MIGHT BE NEEDED LESS REFINED ELSEWHERE."*

YESNObenefits every consumer?current and futuredo it in the pipelinetimestamps, types, lineage, PII removalleave it to the consumer

The test for “should this transformation happen in the pipeline?” is a single question: does it benefit every current and future consumer?

Figure 9.2.9

2.6 Security

The five questions:

  • Who has access to the data that is ingested into Kafka?
  • Can we make sure the data going through the pipe is encrypted? (“mainly a concern for data pipelines that cross datacenter boundaries”)
  • Who is allowed to make modifications to the pipelines?
  • If the pipeline reads/writes from access-controlled locations, can it authenticate properly?
  • Is our PII handling compliant with laws and regulations regarding its storage, access, and use?

What Kafka provides: encryption on the wire (source→Kafka and Kafka→sink), authentication via SASL, authorization ("so you can be sure that if a topic contains sensitive information, it can't be piped into less secured systems by someone unauthorized"), and an audit log to track access — unauthorized and authorized. "With some extra coding, it is also possible to track where the events in each topic came from and who modified them, so you can provide the entire lineage for each record." (Ch. 11.)

⚠️ Connector credentials — do not put them in config files

*"Kafka Connect and its connectors need to be able to connect to, and authenticate with, external data systems, and configuration of connectors will include CREDENTIALS.

These days it is NOT RECOMMENDED to store credentials in configuration files, since this means the configuration files have to be handled with extra care and have restricted access. A common solution is to use an external secret management system such as HashiCorp Vault."*

external secret configurationsupported by Kafka ConnectApache Kafka shipscommunity-developedthe frameworkpluggable external config providersan example providerreads configuration from a fileVault · AWS · Azure
Figure 9.2.11⚠️ Connector credentials — do not put them in config files

2.7 Failure handling

"Assuming that all data will be perfect all the time is DANGEROUS. It is important to plan for failure handling in advance."

The four questions:

  • Can we prevent faulty records from ever making it into the pipeline?
  • Can we recover from records that cannot be parsed?
  • Can bad records get fixed (perhaps by a human) and reprocessed?
  • ⚠ “What if the bad event looks exactly like a normal event and you only discover the problem a few days later?”

Kafka's answer: "Because Kafka can be configured to store all events for long periods of time, it is possible to go back in time and recover from errors when needed. This also allows replaying the events stored in Kafka to the target system if they were lost."

That last question is the one that determines your retention setting: your retention must exceed your bug-discovery latency.

2.8 Coupling and agility — three ways coupling sneaks in

① Ad hoc pipelines. “Some companies end up building a custom pipeline for each pair of applications they want to connect.”

  • Logstash → Elasticsearch
  • Flume → HDFS
  • Oracle GoldenGate: Oracle → HDFS
  • Informatica: MySQL + XML → Oracle
  • … and so on

► “tightly couples the data pipeline to the specific end points and creates a mess of integration points that requires significant effort to deploy, maintain, and monitor. It also means that every new system the company adopts will require building additional pipelines, increasing the cost of adopting new technology, and inhibiting innovation.”

② Loss of metadata. “If the data pipeline doesn't preserve schema metadata and does not allow for schema evolution, you end up tightly coupling the software producing the data at the source and the software that uses it at the destination.”

The concrete failure: “If data flows from Oracle to HDFS and a DBA added a new field in Oracle without preserving schema information… either every app that reads data from HDFS will break, or all the developers will need to upgrade their applications at the same time. Neither option is agile.”

► With schema evolution: “each team can modify their applications at their own pace without worrying that things will break down the line.”

③ Extreme processing. “too much processing ties all the downstream systems to decisions made when building the pipelines about which fields to preserve, how to aggregate data, etc. This often leads to constant changes to the pipeline as requirements of downstream applications change, which isn't agile, efficient, or safe.”

► “The more agile way is to preserve as much of the raw data as possible and allow downstream apps, including Kafka Streams apps, to make their own decisions regarding data processing.”

(② is Ch. 1's schema-registry argument again — the lockstep-deploy problem. It shows up in every chapter because it's the same problem at every layer.)


On this page