Skip to main content

4. User Identity EnrichmentEasy

Core Transformations⏱️ ~10 mins

4. User Identity Enrichment

Enterprise Architecture Context

In production stream-processing architectures (Google Cloud Dataflow / Flink), pipeline stages must handle parallel transformations without data loss, managing schema mutations and aggregations across distributed worker workers.

Problem Statement

### Business Context During CRM ingestion, usernames must be formatted and tagged with a verification label before publishing to customer-facing analytical dashboards. ### Problem Statement Write a function `enrich_users(input_pcoll)` that takes a `PCollection` of username strings (e.g. `"Alice"`) and transforms each into the format: `"Name: <name> [Verified]"`.

Key Learning Objectives

  • Understand distributed Apache Beam execution DAG stages and pipeline lifecycle.
  • Apply idiomatic functional Python transforms using the pipe operator |.
  • Ensure data consistency and idempotency across distributed stream workers.

Sample Data Fixtures

Sample Example 1
Input Stream:
['Alice', 'Bob']
Expected Output:
['Name: Alice [Verified]', 'Name: Bob [Verified]']
Sample Example 2
Input Stream:
['Charlie']
Expected Output:
['Name: Charlie [Verified]']
Topics:#Map#Formatting
Support