cancel
Showing results for 
Search instead for 
Did you mean: 
Community Articles
Dive into a collaborative space where members like YOU can exchange knowledge, tips, and best practices. Join the conversation today and unlock a wealth of collective wisdom to enhance your experience and drive success.
cancel
Showing results for 
Search instead for 
Did you mean: 

A Metadata-Driven Backdated Extraction Pattern for IoT/Sensor Pipelines Using Zerobus and Databricks

aayush_410
New Contributor II

Introduction

IoT devices generate a large amount of data every day. In a perfect scenario, all this data arrives on time and our ingestion pipeline processes it using incremental loads. But in real-world IoT systems, this does not always happen.

A device may lose network connectivity and store readings locally. When the connection is restored, it may send old readings along with the latest data. Similarly, a gateway can upload previously missed data after a network issue. Sometimes, we may also need to reload historical data because of a data-quality or source-system issue. This is where backfill is required.

If a pipeline is designed only for incremental loads, processing historical data can require manual code changes or a separate process. Instead, we can build a metadata-driven ingestion framework that uses the same notebook for both daily loads and historical backfills.

The idea is simple: the notebook receives parameters such as the processing date, entity, and source details. Based on these parameters, it extracts and processes the required data.

For example, consider environmental sensors sending temperature and humidity data through a REST API. The data can be processed through a Databricks medallion architecture:

Raw → Cleansed → Stage → Curated

Using the same framework, we can:

  • Process the latest data.
  • Backfill data for a specific date.
  • Backfill data for a date range.
  • Reprocess data after a source-system issue.

This approach avoids creating separate notebooks for backfill processing and keeps the ingestion logic consistent.

In this Article, we will build a simple metadata-driven backfill framework for an IoT REST API and see how the same processing logic can be used for both incremental and historical data loads. We will also briefly look at how Zerobus Ingest can fit into the real-time ingestion side of the solution.

The Problem: When Date Is Hardcoded

Let's first look at a simple ingestion notebook.

In many pipelines, the current date is used directly to get the data:

 

“

from datetime import datetime

today = datetime.now().date()

readings = call_sensor_api(start=today, end=today)

output_path = f"/mnt/raw/sensors/year={today.year}/month={today.month}/day={today.day}"

readings.write.format("delta").save(output_path)

“

This works fine for a daily job. But what happens if we need to reload data from a previous date?

For example:

"Can we reload the sensor readings from March 8?"

With this design, we need to change the code because the notebook always uses today's date.

This becomes a problem when backfill is required frequently. We may end up creating separate logic for historical data, which makes the pipeline harder to maintain and test.

The solution is simple: make the processing date a parameter instead of hardcoding it in the notebook.

Once the date is a parameter, the same notebook can be used for both today's data and historical data.

Step 1: Create a ProcDate Parameter

The first step is to create a single parameter called ProcDate.

In Databricks, we can create a widget for this:

dbutils.widgets.text("ProcDate", "")

raw_proc_date = dbutils.widgets.get("ProcDate")

If ProcDate is empty, we can use today's date. Otherwise, we use the date provided to the notebook.

from datetime import datetime, date

 

def parse_proc_date(raw_value: str) -> date:

    if raw_value.strip() == "":

        return datetime.utcnow().date()

 

    return datetime.strptime(

        raw_value.strip(),

        "%Y-%m-%d"

    ).date()

 

proc_date = parse_proc_date(raw_proc_date)

Now we have one date variable that controls the entire execution.

For example:

  • ProcDate = 2026-09-04 → process September 4 data.
  • ProcDate = 2026-03-08 → process March 8 data.
  • ProcDate = blank → process today's data.

The important point is to parse the value once and use the same proc_date throughout the notebook. This also helps avoid date and timezone issues caused by handling the date differently in different parts of the code.

aayush_410_0-1789470604258.png

Step 2: Use ProcDate Everywhere

Creating ProcDate is only the first step. The important part is to make sure that all date-related logic uses this value.

There are three main places where we need to use it:

2.1 Output Path

The output path should be created using proc_date.

def build_output_path(

    base_path: str,

    entity: str,

    proc_date: date

) -> str:

 

    return (

        f"{base_path}/{entity}"

        f"/year={proc_date.year}"

        f"/month={proc_date.month:02d}"

        f"/day={proc_date.day:02d}"

    )

 

target_path = build_output_path(

    "/mnt/raw/sensors",

    "temperature_humidity",

    proc_date

)

If we backfill March 8, the data goes to the same partition that a normal run for March 8 would use.

We don't need a separate folder such as /backfill/.

2.2 Store the Date in the Data

We should also store the source date in the data itself.

For example:

from pyspark.sql.functions import lit

 

readings_df = readings_df.withColumn(

    "SourceDate",

    lit(proc_date.isoformat())

)

This is useful because downstream layers can identify the date associated with the reading without depending on the file creation time.

For example, if March 8 data is loaded on March 20, the record should still show March 8 as its SourceDate.

This becomes important for deduplication and other processing in the Cleansed and Curated layers.

2.3 Use ProcDate in the API Request

The API request should also use the same date.

from datetime import datetime, timedelta

 

def build_query_window(proc_date: date):

 

    start = datetime.combine(

        proc_date,

        datetime.min.time()

    )

 

    end = start + timedelta(days=1)

 

    return start, end

 

window_start, window_end = build_query_window(proc_date)

We can then pass this window to the API:

response = requests.get(

    f"{API_BASE_URL}/readings",

    params={

        "start_time": window_start.isoformat(),

        "end_time": window_end.isoformat(),

        "device_group": "temperature_humidity"

    },

    headers={

        "Authorization": f"Bearer {get_api_token()}"

    }

)

Now the API request, output path, and data columns all use the same date.

This gives us one consistent date throughout the pipeline.

aayush_410_1-1789470651275.png

Step 3: Move Entity Configuration to Metadata

So far, we have used one type of sensor. In a real IoT platform, there can be many types of devices and APIs.

For example:

  • Temperature and humidity sensors
  • Door sensors
  • Vibration sensors

Each entity may have a different API endpoint, pagination method, frequency, and target table.

Instead of creating one notebook for every entity, we can store these details in a metadata table.

A simple metadata table could contain:

Entity Name

API Endpoint

Frequency

Pagination

Target Table

temperature_humidity

/v1/readings/temp-humidity

15 min

cursor

raw.temp_humidity_readings

door_contact_sensors

/v1/readings/door-contact

5 min

offset

raw.door_contact_readings

vibration_sensors

/v1/readings/vibration

1 min

cursor

raw.vibration_readings

The notebook can then take EntityName as another parameter:

dbutils.widgets.text("EntityName", "")

 

entity_name = dbutils.widgets.get("EntityName")

It can read the configuration from the metadata table:

entity_config = (

    spark.table("metadata.entity_extraction_config")

    .filter(f"entity_name = '{entity_name}'")

    .collect()[0]

)

 

api_endpoint = entity_config["api_endpoint"]

pagination_style = entity_config["pagination_style"]

target_table = entity_config["target_table"]

Now the same notebook can process different entities.

If we onboard a new sensor type, we mainly need to add its configuration to the metadata table instead of creating another notebook.

aayush_410_2-1789470651310.png

 

Step 4: Handle Backfill-Specific Problems

Once we start using the framework for backfill, a few additional problems can appear.

4.1 Avoid One API Call for Every Day

A simple backfill approach is to run the notebook once for every date:

for single_date in date_range(start_date, end_date):

 

    run_extraction(

        entity_name,

        proc_date=single_date

    )

This can create many API calls.

If the API supports a larger date range, we can process multiple days in one request.

For example, we can create seven-day batches:

def build_batches(

    start_date: date,

    end_date: date,

    batch_days: int = 7

    current = start_date

 

    while current <= end_date:

 

        batch_end = min(

            current + timedelta(days=batch_days - 1),

            end_date

        )

 

        yield current, batch_end

 

        current = batch_end + timedelta(days=1)

Instead of making one request per day, the framework can process the data in larger batches.

aayush_410_3-1789470651336.png

 

4.2 Handle API Rate Limits

Backfill jobs may run for multiple entities at the same time. If these entities use the same API key, the API can return:

429 Too Many Requests

Instead of always waiting for a fixed amount of time, we should check the API's Retry-After header.

def handle_rate_limit(response):

 

    retry_after = response.headers.get("Retry-After")

 

    if retry_after is None:

        return 30

 

    if retry_after.strip().isdigit():

        return int(retry_after)

 

    retry_time = parsedate_to_datetime(retry_after)

 

    wait_seconds = (

        retry_time - datetime.utcnow()

    ).total_seconds()

 

    return max(wait_seconds, 1)

This allows the pipeline to use the wait time provided by the API instead of guessing.

4.3 Avoid Partition Conflicts

Another issue can occur when the pipeline checks for the latest available partition.

During a backfill, the pipeline may accidentally see the partition it is currently writing as the latest partition.

We can avoid this by excluding the current target path:

def get_latest_partition(

    base_path: str,

    exclude_path: str = None

) -> str:

 

    all_partitions = list_partitions(base_path)

 

    if exclude_path:

        all_partitions = [

            p for p in all_partitions

            if p != exclude_path

        ]

 

    return max(

        all_partitions,

        key=extract_date_from_path

    )

Then:

latest = get_latest_partition(

    base_path,

    exclude_path=target_path

)

This prevents the current backfill output from being treated as existing data.

4.4 Avoid Wildcard Imports

One small Python issue can also cause confusing errors.

Avoid:

from pyspark.sql.functions import *

Some PySpark functions such as sum, min, and max can conflict with Python's built-in functions.

It is safer to import only the functions that are required.

For example:

from pyspark.sql.functions import col, lit, max

Keeping imports explicit makes the code easier to understand and reduces the chance of these conflicts.

Step 5: Test the Backfill Framework

Backfill logic is usually not executed as often as the normal daily pipeline. Because of this, it is important to test both normal and historical scenarios.

A simple test plan should include:

Test

Expected Result

ProcDate is blank

Today's date is used

Single-day backfill

Correct date is extracted

Multi-day backfill

All requested dates are processed

API returns 429

Pipeline waits and retries

Live and backfill run together

No partition conflict

Invalid date is provided

Pipeline fails with a clear error

These tests help make sure that changes to the normal ingestion logic do not break the backfill process.

Using Zerobus Ingest for Live IoT Data

The backfill framework above is based on pulling historical data from an API. For live IoT data, we can use a streaming approach.

Zerobus Ingest is a Databricks service for streaming data directly into Delta tables. Instead of a scheduled job repeatedly calling an API, devices or gateways can send records continuously.

This makes Zerobus a good option for the live ingestion path, while the metadata-driven notebook can continue to handle historical backfills.

The main difference is:

  • Live ingestion: data is pushed continuously.
  • Backfill: data is requested for a known historical period.

The important thing is to keep the data contract consistent. If historical records are replayed, they should keep their original SourceDate instead of using the date when the records were replayed.

This allows downstream layers to process live and replayed records in the same way.

Putting It All Together

The complete solution can be summarized in a few steps:

  1. Create a single ProcDate parameter.
  2. Use ProcDate for the API request, output path, and source-date columns.
  3. Store entity-specific configuration in a metadata table.
  4. Use one generic notebook for different entities.
  5. Batch historical API requests when the API supports it.
  6. Handle API rate limits using Retry-After.
  7. Prevent partition conflicts during backfill.
  8. Test both normal and historical execution paths.
  9. Use Zerobus for the live streaming path when it fits the use case.

The result is one ingestion framework that can support both:

Daily ingestion

ProcDate = Today

      ↓

Same Notebook

      ↓

API

      ↓

Delta

and

Historical backfill

ProcDate = Historical Date

      ↓

Same Notebook

      ↓

API

      ↓

Same Delta Structure

No separate backfill notebook is required.

Conclusion

IoT data does not always arrive on time. Devices can lose connectivity, gateways can send old data later, and historical data may need to be reprocessed.

Because of this, backfill should be considered part of the ingestion design rather than something added later.

The main idea in this framework is simple: don't make the pipeline assume the date. Pass the date to the pipeline.

Once ProcDate is used consistently across the API request, output path, and data, the same notebook can process both current and historical data.

Adding metadata makes the framework reusable across different IoT entities, while batching, rate-limit handling, and proper testing make it more reliable for larger backfill workloads.

Aayush Sharma
0 REPLIES 0