Handling a continuous flow of records with Spark's structured streaming
Spark Structured Streaming in Microsoft Fabric enables near-real-time processing of data within notebooks using DataFrames, treating a live data stream as an unbounded table that is incrementally processed. It integrates with Fabric's lakehouse architecture, allowing you to read from sources like files or event streams and write results to Delta tables using familiar Spark APIs.
Must-know
- Structured Streaming uses the same DataFrame/Dataset API as batch processing, so most transformations (select, filter, groupBy, joins) work identically on streaming DataFrames.
- You create a streaming DataFrame with spark.readStream (e.g., specifying format like 'cloudFiles' for Autoloader-style ingestion or 'delta' for reading a Delta table as a stream) and write output with df.writeStream.
- Output modes matter: 'append' only writes new rows and is required for most non-aggregated queries, 'complete' rewrites the full result table (used with aggregations), and 'update' writes only changed rows.
- Checkpointing (via the checkpointLocation option) is mandatory for fault tolerance: it tracks progress and offsets so a restarted stream resumes exactly where it left off without data loss or duplication.
- Trigger intervals control processing cadence: default is micro-batch as fast as possible, but you can set a fixed interval (e.g., trigger(processingTime='1 minute')) or use Trigger.AvailableNow/once for batch-like one-time processing of available data.
- In Fabric, streaming writes commonly target Delta tables in the Lakehouse, enabling downstream consumption via SQL endpoints, Power BI, or further Spark jobs; watermarking is used to handle late-arriving data in stateful aggregations.
A data engineer builds a Fabric notebook that reads from an event stream and uses Spark Structured Streaming to write the results to a Delta table in a lakehouse with writeStream. The job needs to be restartable so that after a failure or a manual stop and restart, it resumes exactly where it left off without reprocessing or losing records. Which configuration should the engineer set on the streaming query?
What you have tried across DP-700's objectives, not a readiness score.
Implement and manage an analytics solution
- Tuning a workspace's Spark compute defaults and pool sizing
- Grouping and governing workspaces with a Fabric domain
- Setting per-workspace defaults for OneLake storage
- Standing up an Airflow job runtime inside a workspace
- Connecting a workspace to a Git repository
- Managing schema changes with a database project
- Promoting Fabric items across environments with a deployment pipeline
- Granting and restricting access at the workspace level
- Locking down who can open a single Fabric item
- Layering row, column, object, and file-level security rules
- Hiding sensitive column values behind a dynamic mask
- Classifying Fabric items with a sensitivity label
- Marking a trusted item as promoted or certified
- Reading a Fabric audit log to see who did what
- Securing data at the OneLake storage layer
- Picking the right build tool among a dataflow, a pipeline, and a notebook
- Kicking off a job on a schedule or in response to an event
- Chaining notebooks and pipelines together with parameters and dynamic expressions
Ingest and transform data
- Deciding between a full reload and an incremental load
- Shaping source data ahead of a dimensional-model load
- Landing a continuous stream of data into storage
- Matching a workload to the right Fabric data store
- Picking a transformation tool from dataflows, notebooks, KQL, or T-SQL
- Linking to external data without copying it via a OneLake shortcut
- Keeping a source database continuously replicated into Fabric
- Moving data into Fabric with a data pipeline
- Writing transform logic in PySpark, SQL, or KQL
- Flattening related tables into one wide, denormalized shape
- Rolling records up with group-by aggregations
- Dealing with duplicate rows, gaps, and data that arrives late
- Selecting the right engine for a real-time workload
- Weighing storage-in-place against a linked shortcut for a Real-Time Intelligence table
- Weighing an accelerated shortcut against a standard one for query speed
- Routing and reshaping live events with an Eventstream
- Handling a continuous flow of records with Spark's structured streaming
- Querying and reshaping event data with KQL
- Aggregating a stream over sliding or tumbling time windows
Monitor and optimize an analytics solution
- Watching an ingestion job's health and progress
- Watching a transformation job's health and progress
- Tracking whether a semantic model's refresh actually succeeded
- Setting up an alert to catch a failure early
- Tracking down why a pipeline run failed and fixing it
- Diagnosing why a dataflow run failed
- Debugging a notebook run that failed
- Troubleshooting a misbehaving Eventhouse
- Troubleshooting a misbehaving Eventstream
- Debugging a T-SQL statement that failed
- Fixing a broken or unreachable shortcut
- Speeding up a Lakehouse table with maintenance operations
- Making a slow pipeline run faster
- Tuning a Fabric warehouse for faster queries
- Improving throughput on real-time streaming components
- Tuning a Spark job to run faster and cheaper
- Making a slow query run faster
Coverage checked against the published exam guide on Aug 11, 2026.
These are independent practice questions, written against this certification's published exam guide. They are not the certification vendor's own questions, and not the real exam.