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.