Skip to content
290 changes: 290 additions & 0 deletions website/docs/quickstart/user_profile.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,290 @@
---
title: Real-Time User Profile
sidebar_position: 4
---

# Real-Time User Profile

This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless.

## How the System Works

### Core Concepts

- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required.
- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine.
- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required.

### Data Flow

1. **Ingestion**: Raw click events arrive with an email address and a click count.
2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically.
3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required.

## Prerequisites

Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine.

## Environment Setup

1. Create a working directory and navigate into it.
```shell
mkdir fluss-user-profile
cd fluss-user-profile
```

2. Create a `docker-compose.yml` file with the following content:
```yaml
services:
coordinator-server:
image: apache/fluss:$FLUSS_DOCKER_VERSION$
command: coordinatorServer
depends_on:
- zookeeper
environment:
- |
FLUSS_PROPERTIES=
zookeeper.address: zookeeper:2181
bind.listeners: FLUSS://coordinator-server:9123
remote.data.dir: /tmp/fluss/remote
volumes:
- fluss-remote-data:/tmp/fluss/remote
tablet-server:
image: apache/fluss:$FLUSS_DOCKER_VERSION$
command: tabletServer
depends_on:
- coordinator-server
environment:
- |
FLUSS_PROPERTIES=
zookeeper.address: zookeeper:2181
bind.listeners: FLUSS://tablet-server:9123
data.dir: /tmp/fluss/data
remote.data.dir: /tmp/fluss/remote
volumes:
- fluss-remote-data:/tmp/fluss/remote
zookeeper:
restart: always
image: zookeeper:3.9.2
jobmanager:
image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$
ports:
- "8083:8081"
command: jobmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
rest.address: jobmanager
rest.port: 8081
volumes:
- fluss-remote-data:/tmp/fluss/remote
taskmanager:
image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$
depends_on:
- jobmanager
command: taskmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
taskmanager.numberOfTaskSlots: 2
volumes:
- fluss-remote-data:/tmp/fluss/remote
sql-client:
image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$
command: ["/opt/sql-client/sql-client"]
depends_on:
- jobmanager
environment:
- |
FLINK_PROPERTIES=
jobmanager.rpc.address: jobmanager
rest.address: jobmanager
rest.port: 8081
volumes:
- fluss-remote-data:/tmp/fluss/remote

volumes:
fluss-remote-data:
```

3. Start all services.
```shell
docker compose up -d
```

4. Confirm all containers are running.
```shell
docker compose ps
```
You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state.

:::note
All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file.
:::

## Enter the SQL Client

Use the following command to enter the Flink SQL Client:

```shell
docker compose run --entrypoint bash sql-client -c "

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we simplify this to docker compose run sql-client, as in the existing Flink quickstart? The service already sets command: ["/opt/sql-client/sql-client"], and the required REST settings are provided through FLINK_PROPERTIES, so the custom entrypoint and duplicated arguments are unnecessary. Since this tutorial creates its own source table, we may also want to point the service command at /opt/flink/bin/sql-client.sh to avoid preloading the unrelated demo tables from sql-client.sql.

\${FLINK_HOME}/bin/sql-client.sh \
-Drest.address=jobmanager \
-Drest.port=8081 \
-i /opt/sql-client/sql/sql-client.sql
"
```

## Step 1: Create the Fluss Catalog

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

All steps should use H3 to make them under the Enter the SQL Client section.


Run these statements one by one in the SQL Client.

:::tip
Run SQL statements one by one to avoid errors.
:::

```sql
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'coordinator-server:9123'
);
```

```sql
USE CATALOG fluss_catalog;
```

:::note
Once you switch to the Fluss catalog, all RoaringBitmap SQL functions (`rb_build_agg`, `rb_cardinality`, `rb_or_agg`, and others) are available immediately — no `CREATE TEMPORARY FUNCTION` statement is needed.
:::

## Step 2: Create the User Dictionary Table

Create the `user_dict` table to map email addresses to integer UIDs. The `auto-increment.fields` property instructs Fluss to automatically assign a unique `INT` UID for every new email it receives.

```sql
CREATE TABLE user_dict (
email STRING,
uid INT,
PRIMARY KEY (email) NOT ENFORCED
) WITH (
'auto-increment.fields' = 'uid'
);
```

## Step 3: Create the Aggregated Profile Table

Create the `user_profiles` table using the **Aggregation Merge Engine**. Each user's UID is the primary key. `total_clicks` is summed and `unique_visitors` accumulates a [RoaringBitmap](https://roaringbitmap.org/) of all UIDs seen — both computed directly at the storage layer.

```sql
CREATE TABLE user_profiles (
uid INT,
total_clicks BIGINT,
unique_visitors BYTES,
PRIMARY KEY (uid) NOT ENFORCED
) WITH (
'table.merge-engine' = 'aggregation',
'fields.total_clicks.agg' = 'sum',
'fields.unique_visitors.agg' = 'rbm32'
);
```

## Step 4: Ingest and Process Data

Create a temporary source table to simulate raw click events using the Faker connector.

:::note
Java Faker's `numberBetween(min, max)` treats `max` as exclusive. The expression below produces click counts of 1–10.
:::

```sql
CREATE TEMPORARY TABLE raw_events (
email STRING,
click_count INT,
proctime AS PROCTIME()
) WITH (
'connector' = 'faker',
'rows-per-second' = '1',
'fields.email.expression' = '#{internet.emailAddress}',
'fields.click_count.expression' = '#{number.numberBetween ''1'',''11''}'
);
```

Now run the pipeline. The `lookup.insert-if-not-exists` hint ensures that if an email is not found in `user_dict`, Fluss generates a new `uid` automatically. `rb_build_agg(d.uid)` builds a one-element RoaringBitmap from each UID — the Aggregation Merge Engine OR-s it into the stored bitmap, giving an exact unique visitor count per user over time.

```sql
INSERT INTO user_profiles
SELECT
d.uid,
CAST(e.click_count AS BIGINT),
rb_build_agg(d.uid)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The user_profiles row key and the bitmap element are both d.uid. For a row keyed by UID u, every incoming bitmap is therefore {u}, so the rbm32 union remains {u} and rb_cardinality(unique_visitors) can never grow beyond 1. This does not demonstrate unique visitors as the tutorial claims. Please introduce a separate aggregation key, such as a page, campaign, or profile-group ID, as the table primary key and keep d.uid as the visitor ID stored in the bitmap.

FROM raw_events AS e

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The uid column from user_dict is never referenced in the SELECT. The lookup join exists solely to trigger insert-if-not-exists, but the generated ID plays no role in the aggregation. This makes the pipeline feel contrived — a reader would expect the dictionary-mapped uid to be the primary key of user_profiles, not an unrelated profile_group_id.

JOIN user_dict /*+ OPTIONS('lookup.insert-if-not-exists' = 'true') */
FOR SYSTEM_TIME AS OF e.proctime AS d
ON e.email = d.email
GROUP BY d.uid, e.click_count;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Using rb_build_agg here introduces an unbounded Flink group aggregation keyed by (uid, click_count). Since uid is high-cardinality and these groups never close, Flink must retain continuously growing keyed state, which defeats the purpose of offloading aggregation to the Fluss Aggregation Merge Engine. We can remove the GROUP BY and use the scalar rb_build(ARRAY[d.uid]) to emit a singleton bitmap for each event; the Fluss rbm32 field will union these bitmaps, while sum aggregates the click counts on the server side.

```

## Step 5: Verify Results

Open a **second terminal**, navigate to the working directory, and launch another SQL Client session to query results while the pipeline runs.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why do we need a second terminal here? Flink SQL Client executes DML asynchronously by default (table.dml-sync=false), so the prompt becomes available again after the INSERT INTO job is submitted. Since this tutorial never enables synchronous DML, we can run the verification queries in the same session and remove the duplicated SQL Client startup and catalog setup.


```shell
docker compose run --entrypoint bash sql-client -c "
\${FLINK_HOME}/bin/sql-client.sh \
-Drest.address=jobmanager \
-Drest.port=8081
"
```

Set up the catalog:

```sql
CREATE CATALOG fluss_catalog WITH (
'type' = 'fluss',
'bootstrap.servers' = 'coordinator-server:9123'
);
USE CATALOG fluss_catalog;
SET 'sql-client.execution.result-mode' = 'tableau';
```

Query the aggregated profile table. `rb_cardinality` converts the stored bitmap into a human-readable unique visitor count:

```sql
SELECT
uid,
total_clicks,
rb_cardinality(unique_visitors) AS unique_visitor_count
FROM user_profiles;
```

You should see rows appearing for each new user, with `total_clicks` and `unique_visitor_count` growing in real time.

To verify the email-to-UID dictionary mapping:

```sql
SELECT * FROM user_dict LIMIT 10;
```

Each email should have a unique compact `INT` uid automatically assigned by Fluss.

## Clean Up

Exit the SQL Client by typing `exit;`, then stop all services.

```shell
docker compose down -v
```

## Architectural Benefits

- **Stateless Flink Jobs:** Offloading identity mapping, click aggregation, and bitmap union to Fluss makes the Flink job lightweight, with fast checkpoints and minimal recovery time.
- **Compact Storage:** Using auto-incremented `INT` UIDs instead of raw email strings reduces memory and storage footprint significantly.
- **Exact Unique Counting:** RoaringBitmap provides exact distinct counts — no approximations like HyperLogLog.
- **Exactly-Once Accuracy:** The Undo Recovery mechanism in the Fluss Flink connector ensures replayed data during failovers does not result in double-counting.

## What's Next?

For the full reference of all RoaringBitmap SQL functions available in FlussCatalog (`rb_or_agg`, `rb_and`, `rb_contains`, `rb_to_array`, and more), see the [SQL Functions](../../engine-flink/sql-functions/) documentation.