You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Commit 6cddc91
Browse filesBrowse the repository at this point in the historyBrowse files
Copy file name to clipboardExpand all lines: docs/modules/demos/pages/data-lakehouse-iceberg-trino-spark.adoc
+61-63Lines changed: 61 additions & 63 deletions
Display the source diff
Display the rich diff
Original file line number
Diff line number
Diff line change
@@ -230,13 +230,11 @@ For details on the NiFi workflow ingesting water-level data, read the xref:nifi-
230
230
231
231
== Spark
232
232
233
-
https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html[Spark Structured Streaming] is used to
234
-
stream data from Kafka into the lakehouse.
233
+
https://spark.apache.org/docs/latest/structured-streaming-programming-guide.html[Spark Structured Streaming] is used to stream data from Kafka into the lakehouse.
235
234
236
235
=== Accessing the web interface
237
236
238
-
To have access to the Spark web interface you need to run the following command to forward port 4040 to your local
239
-
machine.
237
+
To have access to the Spark web interface you need to run the following command to forward port 4040 to your local machine.
The demo has started all the running streaming jobs. Look at the {demo-code}[demo code] to see the actual code
265
-
submitted to Spark. This document will explain one specific ingestion job - `ingest water_level measurements`.
264
+
The demo has started all the running streaming jobs. Look at the {demo-code}[demo code] to see the actual code submitted to Spark.
265
+
This document will explain one specific ingestion job - `ingest water_level measurements`.
266
266
267
-
The streaming job is written in Python using `pyspark`. First off, the schema used to parse the JSON coming from Kafka
268
-
is defined. Nested structures or arrays are supported as well. The schema differs from job to job.
267
+
The streaming job is written in Python using `pyspark`.
268
+
First off, the schema used to parse the JSON coming from Kafka is defined.
269
+
Nested structures or arrays are supported as well.
270
+
The schema differs from job to job.
269
271
270
272
[source,python]
271
273
----
@@ -276,16 +278,13 @@ schema = StructType([ \
276
278
])
277
279
----
278
280
279
-
Afterwards, a streaming read from Kafka is started. It reads from our Kafka at `kafka:9093` with the topic
280
-
`water_levels_measurements`. When starting up, the job will ready all the existing messages in Kafka (read from
281
-
earliest) and will process 50000000 records as a maximum in a single batch. As Kafka has retention set up, Kafka records
282
-
might alter out of the topic before Spark has read the records, which can be the case when the Spark application wasn't
283
-
running or crashed for too long. In the case of this demo, the streaming job should not error out. For a production job,
284
-
`failOnDataLoss` should be set to `true` so that missing data does not go unnoticed - and Kafka offsets need to be
285
-
adjusted manually, as well as some post-loading of data.
281
+
Afterwards, a streaming read from Kafka is started. It reads from our Kafka at `kafka:9093` with the topic `water_levels_measurements`.
282
+
When starting up, the job will ready all the existing messages in Kafka (read from earliest) and will process 50000000 records as a maximum in a single batch.
283
+
As Kafka has retention set up, Kafka records might alter out of the topic before Spark has read the records, which can be the case when the Spark application wasn't running or crashed for too long.
284
+
In the case of this demo, the streaming job should not error out.
285
+
For a production job, `failOnDataLoss` should be set to `true` so that missing data does not go unnoticed - and Kafka offsets need to be adjusted manually, as well as some post-loading of data.
286
286
287
-
*Note:* The following Python snippets belong to a single Python statement but are split into separate blocks for better
288
-
explanation.
287
+
*Note:* The following Python snippets belong to a single Python statement but are split into separate blocks for better explanation.
289
288
290
289
[source,python]
291
290
----
@@ -300,8 +299,7 @@ spark \
300
299
.load() \
301
300
----
302
301
303
-
So far we have a `readStream` reading from Kafka. Records on Kafka are simply a byte-stream, so they must be converted
304
-
to strings and the json needs to be parsed.
302
+
So far we have a `readStream` reading from Kafka. Records on Kafka are simply a byte-stream, so they must be converted to strings and the json needs to be parsed.
305
303
306
304
[source,python]
307
305
----
@@ -319,10 +317,10 @@ Have a look at the {spark-streaming-docs}[Spark streaming documentation on Kafka
First, the data frame containing the upserts (records from Kafka) will be registered as a temporary view so that they
367
-
can be accessed via Spark SQL. Afterwards, the `MERGE INTO` statement adds the new records to the lakehouse table.
363
+
First, the data frame containing the upserts (records from Kafka) will be registered as a temporary view so that they can be accessed via Spark SQL.
364
+
Afterwards, the `MERGE INTO` statement adds the new records to the lakehouse table.
368
365
369
-
The incoming records are first de-duplicated (using `SELECT DISTINCT * FROM waterLevelsMeasurementsUpserts`) so that the
370
-
data from Kafka does not contain duplicates. Afterwards, the - now duplication-free - records get added to the
371
-
`lakehouse.water_levels.measurements`, but *only* if they still need to be present.
366
+
The incoming records are first de-duplicated (using `SELECT DISTINCT * FROM waterLevelsMeasurementsUpserts`) so that the data from Kafka does not contain duplicates.
367
+
Afterwards, the - now duplication-free - records get added to the `lakehouse.water_levels.measurements`, but *only* if they still need to be present.
372
368
373
369
=== The Upsert mechanism
374
370
375
-
The `MERGE INTO` statement can be used for de-duplicating data and updating existing rows in the lakehouse table. The
376
-
`ingest water_level stations` streaming job uses the following `MERGE INTO` statement:
371
+
The `MERGE INTO` statement can be used for de-duplicating data and updating existing rows in the lakehouse table.
372
+
The `ingest water_level stations` streaming job uses the following `MERGE INTO` statement:
377
373
378
374
[source,sql]
379
375
----
@@ -389,25 +385,25 @@ WHEN MATCHED THEN UPDATE SET *
389
385
WHEN NOT MATCHED THEN INSERT *
390
386
----
391
387
392
-
First, the data within a batch is de-deduplicated as well. The record containing the station update with the highest
393
-
Kafka timestamp is the newest and will be used during Upsert.
388
+
First, the data within a batch is de-deduplicated as well.
389
+
The record containing the station update with the highest Kafka timestamp is the newest and will be used during Upsert.
394
390
395
-
If a record for a station (detected by the same `station_uuid`) already exists, its contents will be updated. If the
396
-
station is yet to be discovered, it will be inserted. The `MERGE INTO` also supports updating subsets of fields and more
397
-
complex calculations, e.g. incrementing a counter. For details, have a look at the
398
-
{iceberg-merge-docs}[Iceberg MERGE INTO documentation].
391
+
If a record for a station (detected by the same `station_uuid`) already exists, its contents will be updated.
392
+
If the station is yet to be discovered, it will be inserted.
393
+
The `MERGE INTO` also supports updating subsets of fields and more complex calculations, e.g. incrementing a counter.
394
+
For details, have a look at the {iceberg-merge-docs}[Iceberg MERGE INTO documentation].
399
395
400
396
=== The Delete mechanism
401
397
402
-
The `MERGE INTO` statement can de-duplicate data and update existing lakehouse table rows. For details have a look at
403
-
the {iceberg-merge-docs}[Iceberg MERGE INTO documentation].
398
+
The `MERGE INTO` statement can de-duplicate data and update existing lakehouse table rows.
399
+
For details have a look at the {iceberg-merge-docs}[Iceberg MERGE INTO documentation].
404
400
405
401
=== Table maintenance
406
402
407
403
As mentioned, Iceberg supports out-of-the-box {iceberg-table-maintenance}[table maintenance] such as compaction.
408
404
409
-
This demo executes some maintenance functions in a rudimentary Python loop with timeouts in between. When running in
410
-
production, the maintenance can be scheduled using Kubernetes {k8s-cronjobs}[CronJobs] or {airflow}[Apache Airflow],
405
+
This demo executes some maintenance functions in a rudimentary Python loop with timeouts in between.
406
+
When running in production, the maintenance can be scheduled using Kubernetes {k8s-cronjobs}[CronJobs] or {airflow}[Apache Airflow],
411
407
which the Stackable Data Platform also supports.
412
408
413
409
[source,python]
@@ -439,7 +435,8 @@ while True:
439
435
time.sleep(25 * 60) # Assuming compaction takes 5 min run every 30 minutes
440
436
----
441
437
442
-
The scripts have a dictionary of all the tables to run maintenance on. The following procedures are run:
438
+
The scripts have a dictionary of all the tables to run maintenance on.
On the left, select the database `Trino lakehouse`, the schema `house_sales`, and set `See table schema` to
541
-
`house_sales`.
538
+
On the left, select the database `Trino lakehouse`, the schema `house_sales`, and set `See table schema` to `house_sales`.
542
539
543
540
[IMPORTANT]
544
541
====
545
-
The older screenshot below shows how the table preview would look like. Currently, there is an https://github.com/apache/superset/issues/25307[open issue]
546
-
with previewing trino tables using the Iceberg connector. This doesn't affect the execution the following execution of the SQL statement.
542
+
The older screenshot below shows how the table preview would look like. Currently, there is an https://github.com/apache/superset/issues/25307[open issue] with previewing trino tables using the Iceberg connector.
543
+
This doesn't affect the execution the following execution of the SQL statement.
The data used in this demo is a set of gas sensor measurements*.
44
-
The dataset provides resistance values (r-values hereafter) for each of 14 gas sensors.
45
-
In order to simulate near-real-time ingestion of this data, it is downloaded and batch-inserted into a Timescale table.
46
-
It's then updated in-place retaining the same time offsets but shifting the timestamps such that the notebook code can "move through" the data using windows as if it were being streamed.
43
+
The data used in this demo is a set of gas sensor measurements*.
44
+
The dataset provides resistance values (r-values hereafter) for each of 14 gas sensors.
45
+
In order to simulate near-real-time ingestion of this data, it is downloaded and batch-inserted into a Timescale table.
46
+
It's then updated in-place retaining the same time offsets but shifting the timestamps such that the notebook code can "move through" the data using windows as if it were being streamed.
47
47
The Nifi flow that does this can easily be extended to process other sources of (actually streamed) data.
0 commit comments