Databricks Streaming与AWS S3桶相关技术问题咨询
Hey there! Let's tackle your three questions about Databricks Streaming with AWS S3 in detail:
Absolutely! You can track the time taken for both the read (from S3) and write (to S3) stages of your streaming pipeline, as well as the end-to-end round-trip time. Here are two reliable approaches:
Approach 1: Use StreamingQueryListener for batch-level timing
This listener lets you hook into key events of your streaming query, like when a batch starts processing and when it completes. You can calculate the time taken for each batch's read and write operations by capturing timestamps at these events.
from pyspark.sql.streaming import StreamingQueryListener from datetime import datetime class TimingListener(StreamingQueryListener): def onQueryStarted(self, event): print(f"Query started at: {datetime.now()}") self.batch_start_time = datetime.now() def onQueryProgress(self, event): # Track time taken for the current batch's read and write batch_id = event.progress.batchId read_time = event.progress.durationMs.get("addBatch", 0) # Time to fetch data from S3 write_time = event.progress.durationMs.get("commitOffsets", 0) # Time to write data to S3 total_round_trip = (datetime.now() - self.batch_start_time).total_seconds() * 1000 print(f"Batch {batch_id}: Read time from S3: {read_time}ms, Write time to S3: {write_time}ms, Total round-trip: {total_round_trip:.2f}ms") def onQueryTerminated(self, event): print(f"Query terminated at: {datetime.now()}") # Attach the listener to your Spark session spark.streams.addListener(TimingListener()) # Start your streaming query query = (spark.readStream .format("csv") .option("header", "true") .load("s3://your-input-bucket/path") .writeStream .format("csv") .option("header", "true") .option("checkpointLocation", "s3://your-checkpoint-bucket/path") .start("s3://your-output-bucket/path")) query.awaitTermination()
Approach 2: Custom timing in foreachBatch
If you need more granular control (e.g., tracking time per micro-batch logic), you can add timestamp checks directly in the foreachBatch function:
from datetime import datetime from pyspark.sql.functions import lit def process_batch(df, batch_id): start_time = datetime.now() # Read stage: df is the data fetched from S3 for this batch read_end_time = datetime.now() read_duration = (read_end_time - start_time).total_seconds() * 1000 # Your custom processing logic here processed_df = df.withColumn("batch_id", lit(batch_id)) # Write stage to S3 (processed_df.write .format("csv") .option("header", "true") .mode("append") .save(f"s3://your-output-bucket/path/batch-{batch_id}")) write_end_time = datetime.now() write_duration = (write_end_time - read_end_time).total_seconds() * 1000 total_round_trip = (write_end_time - start_time).total_seconds() * 1000 print(f"Batch {batch_id}: Read duration: {read_duration:.2f}ms, Write duration: {write_duration:.2f}ms, Total round-trip: {total_round_trip:.2f}ms") # Attach to your streaming query query = (spark.readStream .format("csv") .option("header", "true") .load("s3://your-input-bucket/path") .writeStream .foreachBatch(process_batch) .option("checkpointLocation", "s3://your-checkpoint-bucket/path") .start()) query.awaitTermination()
You can use Databricks' dbutils.notebook.run() API, but you need to wrap it inside a foreachBatch function (direct calls in streaming transformations won't work—streaming runs distributed code, and dbutils.notebook.run() is a driver-side operation). Here's how to do it:
Key Notes:
- Pass parameters to the target notebook using the
argumentsparameter (as a dictionary). - Ensure the target notebook accepts these parameters via
dbutils.widgets.get(). - Use
foreachBatchto trigger the notebook call per micro-batch.
Step 1: Prepare the target parameterized notebook
In your existing Python Notebook, define widgets to accept parameters:
# Target Notebook: my_parameterized_notebook dbutils.widgets.text("batch_id", "0", "Batch ID") dbutils.widgets.text("temp_s3_path", "", "Temporary S3 Path for Batch Data") batch_id = dbutils.widgets.get("batch_id") temp_s3_path = dbutils.widgets.get("temp_s3_path") # Your notebook logic here (e.g., read data from temp_s3_path and process) print(f"Processing batch {batch_id} from {temp_s3_path}")
Step 2: Call the notebook from your streaming pipeline
def call_notebook_per_batch(df, batch_id): # Write batch data to a temporary S3 path (to pass to the notebook) temp_path = f"s3://your-temp-bucket/batch-{batch_id}" df.write.format("csv").option("header", "true").mode("overwrite").save(temp_path) # Define parameters to pass to the target notebook notebook_params = { "batch_id": str(batch_id), "temp_s3_path": temp_path } # Run the target notebook result = dbutils.notebook.run( path="/Workspace/Users/your-user@domain.com/my_parameterized_notebook", timeout_seconds=300, # Adjust timeout based on your notebook's runtime arguments=notebook_params ) # Optional: Handle the notebook's output (if it returns a value) print(f"Notebook execution result for batch {batch_id}: {result}") # Attach to your streaming query query = (spark.readStream .format("csv") .option("header", "true") .load("s3://your-input-bucket/path") .writeStream .foreachBatch(call_notebook_per_batch) .option("checkpointLocation", "s3://your-checkpoint-bucket/path") .start()) query.awaitTermination()
Important: If you don't need to pass batch data, you can skip writing to a temp path and just pass metadata parameters directly.
readStream to read CSV files? Spark (and Databricks) automatically adds a special _metadata struct column to DataFrames read from file sources. This column includes the full S3 path of the file each record came from. Here's how to extract it:
Example Code:
from pyspark.sql.functions import col stream_df = (spark.readStream .format("csv") .option("header", "true") .load("s3://your-input-bucket/path")) # Extract the full S3 file path from the _metadata column stream_with_file_path = stream_df.withColumn("source_file_path", col("_metadata.file_path")) # Optional: Extract just the filename if needed stream_with_filename = stream_with_file_path.withColumn("source_filename", col("_metadata.file_name")) # Preview the data (for testing) query = stream_with_filename.writeStream.format("console").start() query.awaitTermination(10)
Key Details:
- The
_metadatacolumn is available in Spark 2.4+ and is enabled by default in Databricks. - Other useful fields in
_metadataincludefile_size(size in bytes) andfile_modification_time(timestamp when the file was last modified). - If you don't see the
_metadatacolumn, ensure you're not using any custom readers that disable it—this is rare in standard Databricks environments.
内容的提问来源于stack exchange,提问作者Hiren Changela

