openframe-adapters-db-influxdb
InfluxDB 2.x time-series database adapter for the OpenFrame Microservice Suite.
Part of the openframe-adapters monorepo. Implements BaseRepository[T]
from openframe-core using influxdb-client's native async client
(InfluxDBClientAsync).
InfluxDB vs. row-based CRUD — read this first
BaseRepository[T] (get/list/create/update/delete) was designed
for row-based stores with a primary key. InfluxDB is a time-series
database: data is immutable points — a measurement name, a tag set, a
field set, and a timestamp — written via line protocol and queried with
Flux, not SQL CRUD. There is no primary key. This adapter maps
BaseRepository onto that model by picking explicit interpretations,
documented honestly here (and in repository.py's module docstring, which
is the canonical source if the two ever drift):
| Method | What it actually does | How well it fits |
|---|---|---|
create(entity) |
Writes one new point tagged with the entity's id. | Good fit. |
get(entity_id) |
Returns the most recent point tagged with entity_id, within a bounded time range (lookback, default -30d). |
No natural fit — InfluxDB has no primary key. "Most recent point for this id" is a chosen interpretation, not the only valid one. |
list(limit, offset) |
limit caps the result count (real). offset skips rows in Python after fetching limit + offset points. |
offset is a compatibility shim, not idiomatic InfluxDB access. Real callers should prefer a time-range query via the raw driver (see below), not an arbitrary row offset into a time-ordered result. |
update(entity) |
Looks up the existing point for the id, then writes a new point with the identical tag set and identical timestamp — InfluxDB's storage engine resolves that as a last-write-wins overwrite of the exact same point. | Not a real update — points are immutable. This is the honest mechanism, not a workaround that pretends otherwise. |
delete(entity_id) |
Checks the point exists, then issues a predicate + time-range delete (delete_api().delete(...)). |
Implementable, but coarser-grained than SQL DELETE … WHERE id = ? — it deletes by predicate across a time range, not a single row. |
If your use case is "look up the latest reading for a sensor" or "write one
reading," this adapter fits well. If it's real time-series analysis —
aggregations, downsampling, multi-field joins across time, continuous
queries — use the raw-driver escape hatch below instead of forcing
BaseRepository to do something it was never shaped for.
Installation
pip install openframe-adapters-db-influxdb
Required env vars:
INFLUXDB_URL=http://localhost:8086
INFLUXDB_TOKEN=your-api-token
INFLUXDB_ORG=your-org
INFLUXDB_BUCKET=your-bucket
Quick start
Raw dict mode
from openframe.adapters.db.influxdb import InfluxDBSettings, InfluxDBRepository
settings = InfluxDBSettings() # reads INFLUXDB_* from env
repo = InfluxDBRepository(settings, measurement="readings", id_tag="sensor_id")
reading = await repo.get("sensor-1") # dict | None — most recent point
readings, total = await repo.list(10, 0) # ([dict, ...], int)
created = await repo.create({"sensor_id": "sensor-1", "temperature": 21.5})
updated = await repo.update({"sensor_id": "sensor-1", "temperature": 22.0})
deleted = await repo.delete("sensor-1") # bool
Typed domain mode
from dataclasses import dataclass
from openframe.adapters.db.influxdb import InfluxDBSettings, InfluxDBRepository
@dataclass
class Reading:
sensor_id: str
temperature: float
class ReadingRepository(InfluxDBRepository[Reading]):
_measurement = "readings"
_id_tag = "sensor_id"
def _record_to_entity(self, record) -> Reading:
return Reading(sensor_id=record["sensor_id"], temperature=record["temperature"])
def _entity_to_point(self, entity: Reading, *, time_override=None):
from influxdb_client.client.write.point import Point
point = Point(self._measurement).tag("sensor_id", entity.sensor_id).field("temperature", entity.temperature)
if time_override is not None:
point = point.time(time_override)
return point
settings = InfluxDBSettings()
repo = ReadingRepository(settings)
reading: Reading | None = await repo.get("sensor-1")
Raw driver access (escape hatch)
Every method above is a compatibility layer over a database that was never
row-shaped. For anything beyond simple single-id lookups/overwrites — real
Flux aggregations, time-range queries, downsampling — use the cached
InfluxDBClientAsync directly, the same "escape hatch for niche features"
pattern this ecosystem's Postgres adapter documents for raw SQL:
client = await repo.client()
tables = await client.query_api().query(
'from(bucket: "my-bucket") |> range(start: -1h) '
'|> filter(fn: (r) => r._measurement == "readings") '
'|> aggregateWindow(every: 5m, fn: mean)',
org=settings.influxdb_org,
)
Wiring into an application
For a real service, wire InfluxDBPlugin (the BasePort-satisfying plugin
class) through ApplicationBootstrap.compose() from openframe-core. This
gives you proper lifecycle management — initialize() / health() /
shutdown() — for free, instead of constructing InfluxDBRepository
directly and managing the client yourself:
from openframe.core.runtime import ApplicationBootstrap
from openframe.core.ports import Capability
from openframe.adapters.db.influxdb import InfluxDBPlugin, InfluxDBSettings
settings = InfluxDBSettings() # reads INFLUXDB_* from env
plugin = InfluxDBPlugin(settings, measurement="readings", id_tag="sensor_id")
async with ApplicationBootstrap.compose(plugin) as app:
repo = app.get(Capability.PERSISTENCE) # -> InfluxDBRepository
reading = await repo.get("sensor-1")
# client is closed automatically on exit (plugin.shutdown() ran)
compose() calls plugin.initialize() on entry and plugin.shutdown() on
exit, so the client is created, health-checked, and torn down without any
manual lifecycle code. Requires openframe-core>=3.3.
Reach for a subclassed ApplicationBootstrap (with a configure() method)
only when you need per-port config=/init_timeout= or conditional
registration order; use app.registry as an escape hatch for anything
neither tier covers. The InfluxDBRepository(settings) construction shown
above under "Quick start" remains valid for tests, scripts, or any context
that doesn't need plugin lifecycle management.
This adapter registers under Capability.PERSISTENCE rather than a
dedicated time-series capability — that enum is closed by design, and
nothing about InfluxDB creates the kind of registry-lookup ambiguity that
would justify adding a member to it. See plugin.py's module docstring.
Resilience — circuit breaking under sustained failure
openframe-core>=3.4 ships openframe.core.resilience.CircuitBreakerProxy
— wrap a repository to short-circuit calls after repeated failures instead
of blocking every caller until operation_timeout during a sustained
outage:
from openframe.core.resilience import CircuitBreakerProxy
from openframe.core.tracing import TracingProxy
repo = CircuitBreakerProxy(
TracingProxy(app.get(Capability.PERSISTENCE).get_repository(), prefix="repository.reading"),
failure_threshold=5,
reset_timeout=30.0,
)
Wrap the traced repository, not the reverse — a short-circuited call never
reaches the adapter, so it shouldn't produce a misleading adapter span. No
adapter code changes are needed to support this — CircuitBreakerProxy
wraps from the outside, exactly like TracingProxy.
Async strategy
This adapter uses influxdb-client's native async client
(InfluxDBClientAsync, backed by aiohttp) — the same native-async shape
as the Postgres/Oracle adapters, not the run_in_executor fallback Cassandra
needs. See connection.py's module docstring for the full research finding
(what was verified against the actually-installed driver, and the real
gotcha found along the way: aiohttp connection-level failures never
become influxdb_client.rest.ApiException, so they're caught explicitly).
Configuration
All settings are read from environment variables.
| Env var | Type | Default | Description |
|---|---|---|---|
INFLUXDB_URL |
str |
required | Server URL |
INFLUXDB_TOKEN |
str |
required | API token |
INFLUXDB_ORG |
str |
required | Organization name or ID |
INFLUXDB_BUCKET |
str |
required | Default bucket name |
CLIENT_TIMEOUT_MS |
int |
10000 |
Per-request timeout (ms) |
LOOKBACK |
str |
-30d |
Flux range(start: ...) window for get()/list() |
CONNECTION_TIMEOUT |
float |
30.0 |
Client creation/connectivity-check timeout (s) |
OPERATION_TIMEOUT |
float |
10.0 |
Per-operation timeout (s) |
MAX_RETRIES |
int |
3 |
Max retry attempts |
Exception hierarchy
All exceptions are AdapterError subclasses from openframe.core.exceptions.
Raw influxdb_client/aiohttp exceptions never escape the adapter.
| Situation | Exception |
|---|---|
Cannot connect (DNS, refused, aiohttp.ClientError) |
AdapterConnectionError |
| Token/org rejected (HTTP 401/403) | AdapterConfigurationError |
| Bad Flux query / missing bucket (HTTP 400/404) | AdapterQueryError |
| InfluxDB server error (HTTP 5xx) | AdapterConnectionError (treated as transient — see repository.py) |
| Operation exceeded timeout | AdapterTimeoutError |
Development
# from the package directory
uv venv .venv && source .venv/bin/activate
uv pip install -e ".[dev]"
python -m pytest tests/ -v
Protocol conformance
from openframe.core.ports import BaseRepository
repo = InfluxDBRepository(settings, measurement="readings", id_tag="sensor_id")
assert isinstance(repo, BaseRepository) # True — structural check
No inheritance from the Protocol is required or used.
License
MIT
Metadata
Release files for openframe-adapters-db-influxdb 0.1.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| openframe_adapters_db_influxdb-0.1.0.tar.gz | 28.4 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| openframe_adapters_db_influxdb-0.1.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 51.6 kB
Release files / openframe_adapters_db_influxdb-0.1.0.tar.gz
| Download URL | openframe_adapters_db_influxdb-0.1.0.tar.gz |
|---|---|
| Size | 28.4 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
99e9a59d6cfcaa53ad6e4db1ee737297cb9c0f54ed17bb960b3fa5bd0e306723
|
|
BLAKE2b-256 checksum How to use checksums |
838ed2843a292738457af39265a438aa9d04a38a70962029383ce3eb3beb06a3
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Release files / openframe_adapters_db_influxdb-0.1.0-py3-none-any.whl
| Download URL | openframe_adapters_db_influxdb-0.1.0-py3-none-any.whl |
|---|---|
| Size | 23.2 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
2d14aa848ec3e01b8dbaa03c589b31e6d25365e2c11d4f8bef2d40ec3523bf0a
|
|
BLAKE2b-256 checksum How to use checksums |
70e7e5330abf08f6e560b64af669014abddc524945e7590952595003a571dccf
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|