データエンジニアは、sensor_id、temperature、timestampの列を持つセンサーの読み取りデータを毎秒受信するストリーミングデータフレームを処理したいと考えています。エンジニアは、データのストリーミング中に、過去5分間の各センサーの平均温度を計算する必要があります。
どのコード実装が要件を満たしていますか?
提供された画像からのオプション:
正解:D
正確な抜粋からの包括的かつ詳細な説明:
正解は D です。これは、適切な時間ベースのウィンドウ集計とウォーターマークを使用するためです。これは、イベント時間データの時間ベースの集計に Spark Structured Streaming で必要なパターンです。
構造化ストリーミングに関する Spark 3.5 ドキュメントより:
イベント時間列にスライディングウィンドウを定義し、groupBをwindow()とともに使用して、それらのウィンドウの集計を計算できます。遅延データを処理するには、withWatermark()を使用して、遅延データの到着許容範囲を指定します。(出典:Structured Streaming Programming Guide) オプションDでは、以下の使用法があります。
パイソン
コピー編集
groupBy("sensor_id", window("timestamp","5分"))
agg(avg("温度").alias("avg_temp"))
各sensor_idについて、5分間のイベント時間ウィンドウにわたって平均温度が計算されることを保証します。ロジックを完成させるために、パイプラインの早い段階でwithWatermark("timestamp", "5 minutes")を使用して、遅いイベントを処理することを前提としています。
他のオプションが間違っている理由の説明:
オプション A は静的 DataFrame またはバッチ クエリに適用され、ストリーミング集計には適していない Window.partitionBy を使用します。
オプション B では時間ウィンドウが適用されないため、5 分間の移動平均は計算されません。
オプション C は集計後に誤って ApplyWithWatermark() を適用し、時間ウィンドウが含まれていないため、必要な時間ベースのグループ化が欠落しています。
したがって、オプション D は、時間ウィンドウ ストリーミング集約を計算するためのすべての要件を満たす唯一のオプションです。