L3 demo: a month of sensor data in PostgreSQL#

This notebook loads the Intel Berkeley Lab sensor data (54 motes, about 2.3 million readings over roughly a month) into PostgreSQL, then answers the engineering questions from the notes with SQL: hourly averages, dropout detection, rolling windows, range checks, and a before/after look at an index with EXPLAIN ANALYZE.

Prerequisites. A running PostgreSQL. The simplest way is the sibling docker-compose.yml:

docker compose up -d          # starts postgres on localhost:5432

The database, user, and password below match that compose file. If you cannot install Docker, the next section starts PostgreSQL inside Colab instead. Install the Python dependencies with uv:

uv run --with pandas,psycopg[binary],sqlalchemy jupyter lab

Run this first on Colab#

Colab starts from its own preinstalled environment rather than this course’s uv environment, so run the cell below before anything else. It installs what this notebook needs and Colab does not already have. Outside Colab it does nothing, so you can run it or skip it.

This demo also needs a PostgreSQL server, which Colab does not provide. Point L3_DSN at a database you can reach, or run this one locally.

# Run this first on Colab. Anywhere else this cell does nothing.
#
# Only genuinely missing packages are installed, so Colab's own versions of
# everything it already ships are left alone.
import importlib.util
import subprocess
import sys

REQUIREMENTS = {
    "pandas": "pandas",
    "psycopg": "psycopg[binary]",
    "sqlalchemy": "sqlalchemy",
}


def _missing(module):
    try:
        return importlib.util.find_spec(module) is None
    except ModuleNotFoundError:  # the parent package is absent
        return True


if "google.colab" in sys.modules:
    need = sorted({pip for mod, pip in REQUIREMENTS.items() if _missing(mod)})
    if need:
        print("installing:", " ".join(need))
        subprocess.run([sys.executable, "-m", "pip", "install", "-q", *need], check=True)
    print("Colab setup done." if need else "Colab: nothing to install.")

No Docker? Start PostgreSQL here instead#

Colab cannot run Docker. It is itself a sandboxed container with no daemon and no privileged access, so docker compose up has nothing to talk to. Docker is not the requirement though: PostgreSQL is, and Colab runs as root with apt, so the server can simply be installed here.

Run the cell below only if you have no PostgreSQL of your own. On a machine where docker compose up -d already worked, skip it; it does nothing outside Colab anyway.

Two things to know before relying on this. A Colab runtime is temporary, so the download and the load happen again every session and disappear when it disconnects. And the index comparison later in this notebook is timed on shared, throttled hardware, so expect the numbers to be mushier than the ones in the notes. Nothing here needs PostgreSQL 16 specifically: the newest feature used is a generated column, which arrived in 12.

# Only if you have no PostgreSQL of your own. Does nothing outside Colab.
import os
import subprocess
import sys

if 'google.colab' in sys.modules:
    def run(*cmd, **kw):
        print('$', ' '.join(cmd))
        return subprocess.run(cmd, **kw)

    run('apt-get', '-y', '-qq', 'update', check=True)
    run('apt-get', '-y', '-qq', 'install', 'postgresql', check=True)
    run('service', 'postgresql', 'start', check=True)

    # CREATE ROLE has no IF NOT EXISTS, hence the DO block, so re-running this
    # cell is harmless. createdb has no equivalent at all, so a second run is
    # simply allowed to fail.
    run('sudo', '-u', 'postgres', 'psql', '-c',
        "DO $$ BEGIN IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname='demo') "
        "THEN CREATE ROLE demo LOGIN PASSWORD 'demo' SUPERUSER; END IF; END $$;",
        check=True)
    run('sudo', '-u', 'postgres', 'createdb', '-O', 'demo', 'sensors')

    # The next cell reads this, so nothing downstream changes.
    os.environ['L3_DSN'] = 'postgresql+psycopg://demo:demo@localhost:5432/sensors'
    print('\nPostgreSQL is up, and L3_DSN points at it.')
else:
    print('Not on Colab, so nothing to do here. Use docker compose up -d.')

Connect#

One SQLAlchemy engine for reading query results into pandas, and one plain psycopg connection for the fast COPY load. The DSN is overridable by environment variable so the same notebook works against a hosted instance.

import os
import sqlalchemy as sa

DSN = os.environ.get(
    'L3_DSN', 'postgresql+psycopg://demo:demo@localhost:5432/sensors')
engine = sa.create_engine(DSN)

with engine.connect() as c:
    print(c.execute(sa.text('select version()')).scalar())

Fetch and parse the raw data#

The raw data.txt has eight columns and no header: date time epoch moteid temperature humidity light voltage. It is cached locally on first run. The rows are not in time order, and some moteid values fall outside the documented 1-54, both of which are normal for a sensor-network export and neither of which is announced.

The delimiter is a single space, and getting that wrong is the trap in this file. A transmission that dropped a channel leaves that field empty, so the row reads ... 42.5178  2.65143 with two spaces in it, and 93,879 of the 2,313,682 rows are like that. sep=' ' sees eight fields on every one of them and turns the empty one into NaN, which is what an absent measurement is. The reflex, sep=r'\s+', collapses the two spaces into one, sees seven fields, and shifts every following value one column to the left, so the voltage lands in light and voltage becomes null.

Nothing raises. on_bad_lines does not fire, because pandas pads a short row with NaN rather than calling it bad; the row count is identical either way. The cell after this one shows the difference the only way you can see it, which is by looking at where the nulls ended up. This is the “real data resists the loader” section of the notes, in the one place it costs you something.

import io, urllib.request, zipfile
from pathlib import Path
import pandas as pd

CACHE = Path('.cache'); CACHE.mkdir(exist_ok=True)
txt = CACHE / 'data.txt'
URL = ('https://raw.githubusercontent.com/'
       'linsea423/Intel_Lab_Data/master/data.zip')
if not txt.exists():
    print('downloading', URL)
    with urllib.request.urlopen(URL) as r:
        buf = r.read()
    with zipfile.ZipFile(io.BytesIO(buf)) as z:
        txt.write_bytes(z.read('data.txt'))

cols = ['date', 'time', 'epoch', 'moteid',
        'temperature', 'humidity', 'light', 'voltage']
# sep=' ' and not r'\s+': an empty field is an absent measurement, and a
# parser that collapses whitespace shifts it away. See the text above.
raw = pd.read_csv(txt, sep=' ', names=cols, header=None, engine='c')
print(f'{len(raw):,} rows, {int(raw.isna().any(axis=1).sum()):,} with a missing field')
raw.head()

The same file, parsed the wrong way#

Run this once, look at the two rows of null counts, and then forget the second parser exists. Both read every row and neither complains.

wrong = pd.read_csv(txt, sep=r'\s+', names=cols, header=None,
                    engine='c', on_bad_lines='skip')
print(f'sep=\' \'    {len(raw):,} rows')
print(raw[cols[4:]].isna().sum().to_string())
print(f'\nsep=r\'\\s+\'  {len(wrong):,} rows')
print(wrong[cols[4:]].isna().sum().to_string())

# Where did the 93,878 missing light readings go? Into `light` as voltages.
shifted = wrong[wrong.light.between(2.0, 3.0) & wrong.voltage.isna()]
print(f'\n{len(shifted):,} rows where \'light\' holds a number in the 2-3 V band'
      f' and voltage is null: {wrong.voltage.max():.3g} V is the highest voltage'
      f' it can now see, against {raw.voltage.max():.4g} V in the correct parse.')
del wrong, shifted

Clean, then reshape to the long/tidy form#

We keep only motes in the documented 1-54 range, build a real timestamptz from the date and time, and melt the four measured channels into the tidy (sensor_id, ts, variable, value) shape the schema uses. This is the wide to long move from the notes, done once in pandas before the data ever reaches the database.

df = raw.dropna(subset=['moteid']).copy()
df['moteid'] = df['moteid'].astype(int)
df = df[df.moteid.between(1, 54)]
df['ts'] = pd.to_datetime(df['date'] + ' ' + df['time'],
                          format='mixed', errors='coerce')
df = df.dropna(subset=['ts'])

long = df.melt(
    id_vars=['moteid', 'ts'],
    value_vars=['temperature', 'humidity', 'light', 'voltage'],
    var_name='variable', value_name='value',
).dropna(subset=['value']).rename(columns={'moteid': 'sensor_id'})
# one row per (sensor, ts, variable): drop the occasional duplicate
long = long.drop_duplicates(['sensor_id', 'ts', 'variable'])
print(f'{len(long):,} tidy readings, '
      f'{long.sensor_id.nunique()} sensors, '
      f"{long.variable.nunique()} variables")
long.head()

Create the schema#

Three tables. sensors holds static per-mote metadata, variables holds the unit and plausible range of each measured quantity, and readings holds the facts, tied to the two dimensions by foreign keys. The foreign key on sensor_id is what makes a reading from a non-existent mote impossible rather than merely unlikely.

One deliberate choice for the index section later: readings uses a surrogate id key and is not indexed on (sensor_id, ts) yet, so we can watch an index change the query plan. In production you would make (sensor_id, ts, variable) the primary key, which enforces one row per reading and gives you the (sensor_id, ts) index below for free.

SCHEMA = '''
DROP TABLE IF EXISTS readings;
DROP TABLE IF EXISTS variables;
DROP TABLE IF EXISTS sensors;

CREATE TABLE sensors (
    sensor_id int PRIMARY KEY,
    x_m double precision,
    y_m double precision
);

CREATE TABLE variables (
    variable text PRIMARY KEY,
    unit text NOT NULL,
    lo double precision,
    hi double precision
);

CREATE TABLE readings (
    id        bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    sensor_id int NOT NULL REFERENCES sensors (sensor_id),
    ts        timestamptz NOT NULL,
    variable  text NOT NULL REFERENCES variables (variable),
    value     double precision
);
'''
with engine.begin() as c:
    c.execute(sa.text(SCHEMA))
print('schema created')

Populate the dimensions#

The motes that actually appear in the data, and one row per measured quantity with its unit and a plausible physical range. Real mote coordinates live in the dataset’s mote_locs.txt; we leave x_m/y_m null here and let the foreign key do its job.

import pandas as pd
# insert only sensor_id; x_m/y_m stay NULL (they live in mote_locs.txt)
sensors = pd.DataFrame({'sensor_id': sorted(long.sensor_id.unique())})
sensors.to_sql('sensors', engine, if_exists='append', index=False)

# The same bands the notes use, so the 0-to-50 C figure there and the
# range query below are talking about one rule rather than two.
variables = pd.DataFrame([
    ('temperature', 'degC', 0.0, 50.0),
    ('humidity', '%RH', 0.0, 100.0),
    ('light', 'lux', 0.0, 2000.0),
    ('voltage', 'V', 2.0, 3.5),
], columns=['variable', 'unit', 'lo', 'hi'])
variables.to_sql('variables', engine, if_exists='append', index=False)
print(sensors.shape[0], 'sensors,', variables.shape[0], 'variables')

Load the readings with COPY#

COPY streams the whole table in one operation, which is the difference between a few seconds and an afternoon of row-by-row INSERTs. We hand it a CSV in memory through psycopg’s copy interface.

import psycopg

raw_dsn = engine.url.render_as_string(hide_password=False).replace(
    'postgresql+psycopg://', 'postgresql://')

buf = io.StringIO()
long[['sensor_id', 'ts', 'variable', 'value']].to_csv(
    buf, index=False, header=False)
buf.seek(0)

copy_sql = ('COPY readings (sensor_id, ts, variable, value) '
            'FROM STDIN WITH (FORMAT csv)')
with psycopg.connect(raw_dsn) as conn:
    with conn.cursor() as cur:
        with cur.copy(copy_sql) as cp:
            for chunk in iter(lambda: buf.read(1 << 20), ''):
                cp.write(chunk)
    conn.commit()

def q(sql):
    '''Run SQL and return the result as a DataFrame.'''
    with engine.connect() as c:
        return pd.read_sql_query(sa.text(sql), c)

print(q('SELECT count(*) AS n FROM readings').n[0], 'readings loaded')

The foreign key is not decorative#

A reading for a mote that is not in sensors is rejected on the spot. This is the integrity guarantee a CSV cannot make.

from sqlalchemy.exc import IntegrityError
try:
    with engine.begin() as c:
        c.execute(sa.text(
            "INSERT INTO readings (sensor_id, ts, variable, value) "
            "VALUES (999, now(), 'temperature', 21.0)"))
    print('inserted (unexpected)')
except IntegrityError as e:
    print('rejected, as it should be:')
    print(' ', str(e.orig).splitlines()[0])

Question 1: hourly average temperature per sensor#

date_trunc buckets the irregular readings into clean hours; GROUP BY does the rest. This is the query the whole session is built around.

q('''
SELECT sensor_id,
       date_trunc('hour', ts) AS hour,
       round(avg(value)::numeric, 2) AS avg_temp
FROM   readings
WHERE  variable = 'temperature'
GROUP  BY sensor_id, hour
ORDER  BY sensor_id, hour
LIMIT  8
''')

Question 2: which motes dropped out?#

A mote that went quiet reported far fewer times than a healthy one. HAVING filters on the aggregate, so this is a one-statement dropout report.

q('''
SELECT sensor_id, count(*) AS n_temp_readings
FROM   readings
WHERE  variable = 'temperature'
GROUP  BY sensor_id
HAVING count(*) < 30000
ORDER  BY n_temp_readings
LIMIT  8
''')

Question 3: gaps in reporting, with a window function#

lag reaches back to a sensor’s previous reading, so ts - lag(ts) is the gap since it last reported. Motes aim for one reading about every 31 seconds, so the large gaps are the dropouts, located exactly in time.

q('''
WITH gaps AS (
  SELECT sensor_id, ts,
         ts - lag(ts) OVER (PARTITION BY sensor_id ORDER BY ts) AS gap
  FROM   readings
  WHERE  variable = 'temperature'
)
SELECT sensor_id, ts, gap
FROM   gaps
WHERE  gap > interval '1 hour'
ORDER  BY gap DESC
LIMIT  8
''')

Question 4: a rolling 1-hour average voltage#

The same window machinery with an aggregate and a time-based frame. Because the sampling is irregular, a RANGE frame measured in time is the honest choice, not a fixed number of rows.

q('''
SELECT sensor_id, ts,
       round(value::numeric, 3) AS voltage,
       round(avg(value) OVER (
         PARTITION BY sensor_id ORDER BY ts
         RANGE BETWEEN interval '1 hour' PRECEDING AND CURRENT ROW
       )::numeric, 3) AS voltage_1h_avg
FROM   readings
WHERE  variable = 'voltage' AND sensor_id = 1
ORDER  BY ts
LIMIT  8
''')

Question 5: the impossible readings, and what predicts them#

Roughly a fifth of the temperature readings are physically impossible for an indoor lab. Joining each temperature reading to the same mote’s voltage at the same instant shows why: nearly every impossible temperature comes from a mote whose battery had already fallen below about 2.4 V. Voltage is a data-quality signal, not just housekeeping.

q('''
WITH t AS (
  SELECT sensor_id, ts, value AS temp FROM readings
  WHERE variable = 'temperature'
), v AS (
  SELECT sensor_id, ts, value AS volt FROM readings
  WHERE variable = 'voltage'
)
SELECT
  count(*) FILTER (WHERE temp < 0 OR temp > 50) AS impossible,
  round(100.0 * avg((temp < 0 OR temp > 50)::int), 1) AS pct_impossible,
  round(100.0 * (count(*) FILTER (WHERE (temp < 0 OR temp > 50)
                                    AND volt < 2.4))
        / nullif(count(*) FILTER (WHERE temp < 0 OR temp > 50), 0), 1)
        AS pct_of_impossible_below_2v4
FROM t JOIN v USING (sensor_id, ts)
''')

Indexing: the same range query, before and after#

The query that dominates time-series work is a single sensor over a window of time. Read the plan from the scan at the base (Seq Scan vs Index Only Scan) and the total execution time at the bottom. With no index on (sensor_id, ts), it is a sequential scan of the whole table, parallelized across workers but still touching every row.

def explain(sql):
    for row in q('EXPLAIN ANALYZE ' + sql)['QUERY PLAN']:
        print(row)

RANGE_Q = '''
SELECT count(*) FROM readings
WHERE  sensor_id = 1
  AND  ts BETWEEN '2004-03-15' AND '2004-03-16'
'''
explain(RANGE_Q)

Now add the composite B-tree on (sensor_id, ts) and run the identical query. The top node changes to an index scan and the execution time collapses.

with engine.begin() as c:
    c.execute(sa.text(
        'CREATE INDEX IF NOT EXISTS readings_sensor_ts '
        'ON readings (sensor_id, ts)'))
    c.execute(sa.text('ANALYZE readings'))
explain(RANGE_Q)

Where this goes next#

Everything here is OLTP: correct writes, foreign keys, and fast point and range reads over a row store. L4 keeps this exact dataset and changes only where it lives, into columnar Parquet and the embedded engine DuckDB, and asks when an analytical column store beats the database we just built.