confluentinc / confluentinc/kafka-tutorials

Aviation Tutorial does not compute correct result

Open
#1,289 0 comments 1 reaction 0 assignees View on GitHub
bug recipe
Dominant language
Java
Stars
39
Forks
91
PR merge metrics
No merged PRs in 30d

Description

https://developer.confluent.io/tutorials/aviation/confluent.html

Looking into this recipe, it seems that it would actually not compute the right thing?

We start with 3 tables, `customers`, `flights`, and `bookings`.

In the first two queries, we "enrich" booking with customer and flight details (`booking` is just containing FKs to customer and flights).

In the end, we get `customer_flights` tables which contains a row per booking (booking.id is the PK of this table).

This also means, that there can be two rows with the same `flight.id` if two customer book the same flight.

When we now rekey the data on `flight.id` we actually get "garbage", as it's unclear which of the potentially many bookings/row with the same flight.id will be in customer_flights_rekeyed -- we can only store one row per flight and one row will end up there non-deterministically.

The reason the recipe is doing the re-keying is because of the stream-table join they want to do in the last step: for each flight update, they want to generate some alert for the customer -- however, as only one "random" customer will end up in the table we join to, we can only alert a single customer, but not multiple. Hence, the recipe breaks (besides the fact that we also cannot predict which customer will get the notification).

I think the right thing to do would be to actually not re-key the `customer_flight` table, but to aggregate it via `GROUP BY flight.id` and use `COLLECT_SET` or similar to get all customers with the same `flight.id` in a single row. -- After the stream-table join, you would need to use `EXPLODE` to split a single enriched flight update event, into one alert event per customer.

Thoughts?

(Btw: it might actually be simpler, to just do the `GROUP_BY flight_id` on bookings and just collect all `customer.id`. Do the stream-table joins with the result table of the group-by; split/explode the alerts in the enriched flight-update stream to get an alert-event per customer (which would at this point only contain the customer id), and enrich the alert stream with a second stream-table join to the customer table?)

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.