· Case Study · 16 min read
From Kafka on-prem to a fully managed data platform on Google Cloud
Raw events kept forever, everything else a version — change the code, trigger one DAG, and the history is replayed while live events keep flowing ; more than three billion events a year, and nothing to operate ⏱️

We are moving to Google Cloud, and I want as many managed services as possible. My ops team is five people, for ten to twenty websites in production.
That was the brief on my first day, in January 2017. It came with a story.
The events of the job boards lived on a Kafka cluster in a datacenter. One day a topic was deleted in production, by mistake, and everything it held was almost lost. It came back because a data engineer went to the datacenter and plugged a USB drive into the Kafka servers 💀
That incident is where the brief came from, and it is where the architecture came from. Five years later the platform takes more than three billion events a year through Pub/Sub, Dataflow, BigQuery and Cloud Composer, and nobody on the team operates a broker, a scheduler or a cluster node. It rests on one rule :
The raw events are the only thing that runs forever. Everything else is a version.
The raw side collects every event and keeps it, untouched, 24/7. The enrichment — its code, its tables, its fields, its search index — can be changed whenever the product needs it. We do not migrate data when that happens. We trigger one DAG : 44 tasks that deploy the new jobs, replay the history from the raw events, keep the live events flowing through both versions, and switch when the new one has caught up.
Why we left the datacenter
A Kafka cluster is not a piece of software, it is a team. Brokers to patch, ZooKeeper to keep quorate, disks to size, partitions to rebalance, upgrades to schedule, and a pager for the night a broker fills up. That is a fine investment for a company whose product is the streaming platform. Ours is job boards, and five people who already run ten to twenty websites cannot also be a streaming team.
So the brief settled the stack before any benchmark : only services we do not have to operate. Pub/Sub instead of Kafka, Dataflow instead of a Flink or Spark cluster, BigQuery instead of a warehouse we would have to size — and, as soon as it was an option for us, Cloud Composer instead of the Airflow we hosted ourselves. For a while a small Kafka-to-Pub/Sub bridge ran as a pod, so that the producers in the datacenter could keep publishing as before while the cloud side was built against real traffic. Then they moved to an HTTP tracker that publishes straight to Pub/Sub, and the bridge was switched off.
It also had to be built by very few people : in 2017 two of us committed to the platform’s repositories, and nine commits out of ten were mine. A stack with nothing to operate is what made the rest of this post affordable.
The architecture

Three zones, and they do not have the same life :
- Collect and keep runs 24/7. A tracker receives the events and publishes them, a raw job writes them to BigQuery as they are. Adding a field to an event does not redeploy any of it.
- Derive is a version. The enrichment job parses the events, enriches them, and writes an enriched table and a search index that the services read through an alias. When we change it, we do not patch it : we build the next one.
- Replay exists only when we ask for it. It turns the raw table back into a stream, so that the next version can process five minutes ago and two years ago with the same code.
Raw : four columns, all day, every day
The raw job is the least clever job of the platform, and that is its whole design :
@BigQueryType.toTablecase class Row( event_id: Option[String], event_timestamp: Option[Instant], event_processed_at: Option[Instant], event_raw: Option[String] // the event, as it arrived)
def parseJsonStep(wv: WindowedValue[String]): WindowedValue[Row] = { val json = Json.parse(wv.value) wv.copy(value = Row( event_id = (json \ "id").asOpt[String], event_timestamp = Some(wv.timestamp), event_processed_at = Some(new Instant(System.currentTimeMillis)), event_raw = Some(wv.value) ))}Four columns. An id, two timestamps, and the payload as a string. The job reads one field of the event and interprets none of the others, so it has no opinion about what an event should contain — and no reason to change when the events do. The enriched table next to it has more than fifty columns, and every one of them is a decision somebody may want to revisit.
It was also the first job to be tested for real. When the history of the Kafka cluster was fed into the raw table — the same cluster that had nearly lost a topic — the Dataflow console looked like this :

305,624 events per second through the parser and into BigQuery, and nobody sized a cluster for it. From that day the events had a second home, partitioned by day, that no single command on a broker could take away.
The table is partitioned by day, from the event’s own timestamp :
.to((value: ValueInSingleWindow[TableRow]) => { val day = DateTimeFormat.forPattern("yyyyMMdd").withZone(UTC) .print(value.getTimestamp.getMillis) new TableDestination(s"$project:$dataset.$table$$$day", null)})A day is the unit of everything that follows. A replay is a range of days. A check is a count per day. BigQuery bills the bytes it scans, and a query on one partition scans one day, not the history. Nothing clever — the cheapest unit of undo the warehouse offers, chosen on day one.
Replay : the raw table, as a stream again
The replay job is the mirror of the raw job : a date range in, the same events out.
val sqlQuery = s"SELECT *, TIMESTAMP_TO_MSEC(event_timestamp) as event_timestamp_unix " + s"FROM [$project:$dataset.$table] " + s"WHERE _PARTITIONTIME BETWEEN TIMESTAMP('${opts.getStartDate}') " + s"AND TIMESTAMP('${opts.getEndDate}')"
sc .bigQuerySelect(sqlQuery) .withFixedWindows(Duration.standardSeconds(WINDOW_SIZE)) .toWindowed .map { wv => val row = wv.value val attributes = new util.HashMap[String, String]() // the same "id" and "timestamp" the live events carry attributes.put("id", row.get("event_id").asInstanceOf[String]) attributes.put("timestamp", row.get("event_timestamp_unix").asInstanceOf[String])
val raw = row.get("event_raw").asInstanceOf[String] wv.withValue(new PubsubMessage(encode(raw), attributes)) } .toSCollection .saveAsCustomOutput("SaveToPubSub", PubsubIO.writeMessages().to(topic))The two attributes are the point. id is what deduplicates, timestamp is what makes a window follow the time of the event rather than the time it arrives. With both set, whatever reads the replay topic cannot tell a replayed event from a live one — so there is no “replay mode” to write in any job downstream. The same code, on another subscription.
Enrichment : the part that is allowed to change
The enrichment job is where the product lives : it parses the recruiters’ events into typed records, resolves the user agent, and writes two outputs — a partitioned table for analysis, an index for the statistics the recruiters see. It is also where change requests land. Between October 2019 and July 2021, fourteen commits added fields, columns or whole event types to that model : emails, a transaction id and its origin, a correlation id, the time spent on a résumé, extended search criteria, the label, country and postal code of a location, saved-search events.
Each of them is true for the events of tomorrow the minute it is deployed. It is true for the events of last year only if last year is enriched again.
On the job’s side that takes two flags and one convention :
val rows = sc .withName("ReadPubSub") .read(PubsubIO.string(subscription, "id", "timestamp"))( PubsubIO.ReadParam(isSubscription = true))
// write to BQ if not disabledif (!bigQueryOutputDisabled) { /* day partitions, or dated tables for a replay */ }
// write to ESif (!esOutputDisabled) { esRows.saveAsCustomOutput(s"SaveToEs-$esIndexPrefix", write)}The convention is that an event keeps its id everywhere. It is the deduplication attribute on Pub/Sub, the document id in the index — so an event that arrives once live and once replayed is one document — and the column the tables are deduplicated on.
One trigger, 44 tasks, four jobs
The DAG has no schedule. It is created paused, and it runs when somebody who has just changed the enrichment triggers it in Composer. Then it does, in order, what we used to do with a checklist :

- A copy of the current version takes over the live stream, on its own subscription, writing to the index the services read. Its BigQuery output is switched off : the tables must not be written twice.
- The next version starts next to it, on a new subscription, writing to a new index. From that minute, every live event is enriched by both versions.
- The history is replayed : the replay job reads the raw table and publishes to the replay topic, and a second instance of the next version consumes it into the same new index. Today that is every event since January 2020, and the DAG gives it twenty hours.
- Then the DAG closes the books : wait for the streaming buffer to empty, turn the replay’s day tables into partitions, append, deduplicate ; tune the new index for search, force-merge it ; move the alias ; stop the old job, delete its subscription and its indexes ; write the new version number in a bucket.
Everything that makes this safe is in the names. The version is a number in a bucket, and it is a suffix on the subscription, the index and the job. This is how the DAG launched the next version when I wrote it in 2019, condensed :
launch-job -main acme.beam.jobs.EnrichEvents \ --runner=DataflowRunner \ --streaming=true \ --pubsubTopic={{params.dataflow.pubSubTopic}} \ --pubsubSubscription={{params.dataflow.pubSubSubscription}}-v$((VERSION + 1)) \ --bigQueryTable={{params.bigQuery.tableName}} \ --usePartitioningPerDay=true \ --esIndexPrefix={{params.elasticsearch.indexPrefix}}-v$((VERSION + 1)) \ --jobName={{params.dataflow.jobName}}-v$((VERSION + 1))The replay instance is the same command with the replay subscription, --useDatedTables=true and a table named _replay. And the very last task of the DAG, once everything else has succeeded, is the only place where the number changes :
echo "$((VERSION + 1))" | gsutil cp - gs://${BUCKET}/enrich-events/VERSIONA run that fails halfway has not touched the index the services read. The next run starts by cleaning the job, the subscription and the index the failed one left under the same N+1 names, and goes again.
The last step a service ever notices is one call :
curl -XPOST -H "Authorization: Basic ${BASIC_AUTH}" \ -H 'Content-Type: application/json' \ ${ES_HOST}:${ES_PORT}/_aliases -d '{ "actions" : [ { "add" : { "index" : "'${ES_INDEX_PREFIX}'-v'${VERSION}'-*", "alias" : "'${ES_INDEX_PREFIX}'" } } ]}'Changing the enrichment is a deployment, not a migration.
How it got there
Replay on demand is as old as the platform ; the DAG is not.
- 2017 — the replay job exists in March. In September a replay is two pipelines in the CI : one starts the replay job, the other starts a second instance of the counters job on a
replay-topic. On demand already, by hand, and only for the brave. - 2019 — the enrichment job, and the DAG : thirty-one tasks and three jobs, on the Airflow we hosted ourselves. Each task was a bash operator that pulled the job’s package from a bucket and ran one of the scripts shipped inside it — clean, apply the index templates, wait for the job, check that the replay is finished, dated tables to partitions, wait for the streaming buffer, deduplicate. I wrote them in that order, over four weeks, one script each time the DAG reached a new step.
- 2020 — a teammate ports the DAG to Cloud Composer, and the bash operators become pods running the job’s image.
- 2021 — the fourth job. Until March a rebuild still had a hole in it : the live enrichment of the index being served stopped while the new one was built.
init_enrich_previous_live = { "runner": "DataflowRunner", "streaming": "true", "pubsubSubscription": "{{params.previousLive.pubSubSubscription}}" "-v{{ ti.xcom_pull(task_ids='dump-current-version') }}", "esIndexPrefix": "{{params.elasticsearch.indexPrefix}}" "-v{{ ti.xcom_pull(task_ids='dump-current-version') }}", # the index being served stays fresh, the tables are not written twice "bigQueryOutputDisabled": "true", "jobName": "{{params.previousLive.jobName}}" "-v{{ ti.xcom_pull(task_ids='dump-current-version') }}",}- 2022 — the indexes move to Elasticsearch 7 in March, through the same DAG : a new version, a replay, an alias. The same month, the scheduler moves to Composer 2.
Three schedulers in five years, and the DAG survived each move for one reason : it never imports a job. It launches the job’s container and waits.
launch_enrich = GKEStartPodOperator( task_id='launch-enrich-dataflow', name='launch-enrich-dataflow', namespace=params['podOperator']['namespace'], cluster_name=params['podOperator']['clusterName'], # the Beam job ships as a Docker image ; the DAG only knows its tag image='gcr.io/%s/beam-enrich-events:%s' % (params['project'], params['docker-tag']), # --runner=DataflowRunner, --project, --subscription, ... arguments=launch_args, is_delete_operator_pod=True, secrets=[service_account], dag=dag,)
The same move, one level up : 62 days of dual writes
In 2020 the platform itself had to move. It had outgrown the single project it was born in — three products, one project, one blast radius — and my teammates rebuilt the foundation as code : one Terraform repository per product, a staging and a production project from the same modules, 368 resources today, applied by Cloud Build. Every producer and every consumer still pointed at the old project.
We did with the platform what the DAG does with an index : make the new path live, run both, compare, switch, delete. For the events API, the Go service that publishes the recruiters’ events, it took five steps over the autumn :
- September 2 — the new project learns that the old topics exist.
- September 29 — the API starts publishing to both the old and the new topic.
- October 2 — the old project is allowed to publish into the new one, so nothing upstream is stranded.
- October 6 — the consumers switch to the new topic.
- November 30 — the old topic is removed, 62 days after the first dual write.
type PubSubService struct { client *pubsub.Client topic *pubsub.Topic legacyTopic *pubsub.Topic // FIXME: to remove}
func (s PubSubService) Publish(msg []byte, attrs map[string]string) { result := s.topic.Publish(ctx, message(msg, attrs)) legacyResult := s.legacyTopic.Publish(ctx, message(msg, attrs))
go logOutcome("new", result) go logOutcome("legacy", legacyResult) // FIXME: to remove}Two FIXME: to remove comments, both honoured. And because the raw table exists on each side, “compare” is not a meeting, it is a script from 2017 — one count per day, old against new. Condensed :
while [[ ${CURRENT_PARTITION} -lt ${TODAY} ]]; do
COUNT_1=$(bq query -q --headless --nouse_cache --use_legacy_sql \ --format=json --project_id ${P1} \ "SELECT COUNT(*) as count FROM ${D1}.${T1}\$${CURRENT_PARTITION}" \ | jq -r '.[] | .count')
COUNT_2=$(bq query -q --headless --nouse_cache --use_legacy_sql \ --format=json --project_id ${P2} \ "SELECT COUNT(*) as count FROM ${D2}.${T2}\$${CURRENT_PARTITION}" \ | jq -r '.[] | .count')
echo "${CURRENT_PARTITION} | ${COUNT_1} | ${COUNT_2} |"
CURRENT_PARTITION=$(date -v+1d -j -f '%Y%m%d' \ "${CURRENT_PARTITION}" +'%Y%m%d')doneWhat it gives us today
- Nothing to operate. No broker, no ZooKeeper, no scheduler host, no worker pool. The pager of 2017 is a bill.
- A raw side nobody has to touch. Four columns, one partition per day, 24/7.
- An enrichment anybody can change. A new field is a commit and a trigger. The history follows.
- No cut for the people who read it. Live events go through both versions during a rebuild, and the services change index on one alias.
- An undo. Until the alias moves, the previous version is still there, still fresh.
One limit, because it is true : the index is versioned, the enriched table is not. The DAG deletes it at the start and fills it again, from the live events of the new version and from the replay. Whoever queries that table during a rebuild sees it grow back. The day somebody depends on it the way the services depend on the index, it will need its own alias.
Whose platform it is
In 2017 I wrote nine commits out of ten on the platform’s repositories. This year, fewer than two out of ten.

The DAG I wrote in 2019 is not the one that runs today : a teammate ported it to Composer, and both DAG repositories were started by that same teammate. The résumé index has its own DAG, written by the team that runs it, on the same idea — a current version, a next version, a switch. The biggest Beam job of the platform is theirs too. And the Terraform that made the 2020 move possible is theirs from the first line.
The jobs I wrote prove the platform works. The jobs I did not write prove the design does ✅
Lessons learned
- Keep the raw side dumb. A job that does not interpret the events has no reason to change when they do, and it is the one job you cannot afford to get wrong.
- Make everything else a version. A suffix on the subscription, the index and the job ; one alias ; one number that changes last.
- Replay is a feature, not an incident procedure. If it takes courage, it will be avoided, and the model will stop evolving. If it takes a trigger, it will be used.
- Never stop the live stream to fix the past. Both versions enrich the present while the history catches up.
- Only run what you cannot buy. Every managed service in this stack replaced a pager, and one of them replaced a USB drive.
The other half of this move is the services around the platform — and they are written in Go, because a Scala hello world once refused to fit its pod. That one is for another post 😉



