Skip to content

Repository files navigation

arc-tsdb-client

Python SDK for Arc time-series database.

PyPI versionPython 3.9+License: MIT

Installation

pip install arc-tsdb-client
# With pandas support
pip install arc-tsdb-client[pandas]
# With polars support
pip install arc-tsdb-client[polars]
# With all optional dependencies
pip install arc-tsdb-client[all]

Or with uv:

uv add arc-tsdb-client
uv add arc-tsdb-client --extra pandas
uv add arc-tsdb-client --extra all

Quick Start

fromarc_clientimportArcClientwithArcClient(host="localhost", token="your-token") asclient:
# Write data (columnar format - fastest)client.write.write_columnar(
measurement="cpu",
columns={
"time": [1633024800000000, 1633024801000000],
"host": ["server01", "server01"],
"usage_idle": [95.0, 94.5],
},
)
# Query to pandas DataFramedf=client.query.query_pandas("SELECT * FROM default.cpu WHERE host = 'server01'")
print(df)

Features

  • High-performance ingestion: MessagePack columnar format (9M+ records/sec)
  • Multiple formats: MessagePack columnar/row, InfluxDB Line Protocol
  • DataFrame integration: Pandas, Polars, PyArrow
  • Sync and async APIs: Full async support with httpx
  • Buffered writes: Automatic batching with size and time thresholds
  • Query support: SQL queries with JSON or Arrow IPC streaming
  • Management: Retention policies, continuous queries, delete operations
  • Authentication: Token management (create, rotate, revoke)

Data Ingestion

Columnar Format (Recommended)

The fastest way to write data. Uses MessagePack with gzip compression:

client.write.write_columnar(
measurement="cpu",
columns={
"time": [1633024800000000, 1633024801000000], # microseconds"host": ["server01", "server02"],
"region": ["us-east", "us-west"],
"usage_idle": [95.0, 87.3],
"usage_system": [3.2, 8.1],
},
database="default", # optional
)

DataFrame Ingestion

Write directly from pandas or polars DataFrames:

importpandasaspddf=pd.DataFrame({
"time": pd.to_datetime(["2024-01-01", "2024-01-02"]),
"host": ["server01", "server02"],
"value": [42.0, 43.5],
})
client.write.write_dataframe(
df,
measurement="metrics",
time_column="time",
tag_columns=["host"],
)

Buffered Writes

For high-throughput scenarios, use buffered writes with automatic batching:

withclient.write.buffered(batch_size=10000, flush_interval=5.0) asbuffer:
forrecordinrecords:
buffer.write(
measurement="events",
tags={"source": record.source},
fields={"value": record.value},
timestamp=record.timestamp,
)
# Auto-flushes on exit or when batch_size/flush_interval reached

Line Protocol

For compatibility with InfluxDB tooling:

# Single lineclient.write.write_line_protocol("cpu,host=server01 usage=45.2 1633024800000000000")
# Multiple lineslines= [
"cpu,host=server01 usage=45.2",
"cpu,host=server02 usage=67.8",
]
client.write.write_line_protocol(lines)

Querying Data

JSON Response

result=client.query.query("SELECT * FROM default.cpu WHERE time > now() - INTERVAL '1 hour'")
print(result.columns) # ['time', 'host', 'usage']print(result.data) # [[1633024800000000, 'server01', 45.2], ...]print(result.row_count)

pandas DataFrame

df=client.query.query_pandas("SELECT * FROM default.cpu LIMIT 1000")

Polars DataFrame

pl_df=client.query.query_polars("SELECT * FROM default.cpu LIMIT 1000")

PyArrow Table (Zero-Copy)

table=client.query.query_arrow("SELECT * FROM default.cpu LIMIT 1000")

Query Estimation

Preview query cost before execution:

estimate=client.query.estimate("SELECT * FROM default.cpu")
print(f"Estimated rows: {estimate.estimated_rows}")
print(f"Warning level: {estimate.warning_level}") # none, low, medium, high

List Measurements

measurements=client.query.list_measurements(database="default")
forminmeasurements:
print(f"{m.measurement}: {m.file_count} files, {m.total_size_mb:.1f} MB")

Async Support

All operations have async equivalents:

importasynciofromarc_clientimportAsyncArcClientasyncdefmain():
asyncwithAsyncArcClient(host="localhost", token="your-token") asclient:
# Writeawaitclient.write.write_columnar(
measurement="cpu",
columns={"time": [...], "host": [...], "usage": [...]},
)
# Querydf=awaitclient.query.query_pandas("SELECT * FROM default.cpu")
# Buffered writesasyncwithclient.write.buffered(batch_size=5000) asbuffer:
forrecordinrecords:
awaitbuffer.write(measurement="events", fields={"v": record.v})
asyncio.run(main())

Management Operations

Retention Policies

Automatically delete old data:

# Create policypolicy=client.retention.create(
name="30-day-retention",
database="default",
retention_days=30,
measurement="logs", # optional, applies to all if not specified
)
# List policiespolicies=client.retention.list()
# Execute (with dry_run first)result=client.retention.execute(policy.id, dry_run=True)
print(f"Would delete {result.deleted_count} rows")
# Execute for realresult=client.retention.execute(policy.id, dry_run=False, confirm=True)

Continuous Queries

Aggregate and downsample data automatically:

# Create a CQ to downsample CPU metrics to 1-hour averagescq=client.continuous_queries.create(
name="cpu-hourly",
database="default",
source_measurement="cpu",
destination_measurement="cpu_1h",
query="SELECT time_bucket('1 hour', time) as time, host, avg(usage) as usage FROM default.cpu GROUP BY 1, 2",
interval="1h",
)
# Manual execution with time rangeresult=client.continuous_queries.execute(
cq.id,
start_time="2024-01-01T00:00:00Z",
end_time="2024-01-02T00:00:00Z",
dry_run=True,
)
# List CQscqs=client.continuous_queries.list(database="default")

Delete Operations

Delete data with WHERE clause:

# Always dry_run firstresult=client.delete.delete(
database="default",
measurement="logs",
where="time < '2024-01-01'",
dry_run=True,
)
print(f"Would delete {result.deleted_count} rows from {result.affected_files} files")
# Execute deletionresult=client.delete.delete(
database="default",
measurement="logs",
where="time < '2024-01-01'",
dry_run=False,
confirm=True, # Required for large deletes
)

Authentication

# Verify current tokenverify=client.auth.verify()
ifverify.valid:
print(f"Token: {verify.token_info.name}")
print(f"Permissions: {verify.permissions}")
# Create new tokenresult=client.auth.create_token(
name="my-app-token",
description="Token for my application",
permissions=["read", "write"],
)
print(f"New token: {result.token}") # Save this - shown only once!# List tokenstokens=client.auth.list_tokens()
# Rotate tokenrotated=client.auth.rotate_token(token_id=123)
print(f"New token: {rotated.new_token}")
# Revoke tokenclient.auth.revoke_token(token_id=123)

Configuration

client=ArcClient(
host="localhost", # Arc server hostnameport=8000, # Arc server port (default: 8000)token="your-token", # API tokendatabase="default", # Default databasetimeout=30.0, # Request timeout in secondscompression=True, # Enable gzip compression for writesssl=False, # Use HTTPSverify_ssl=True, # Verify SSL certificates
)

Error Handling

fromarc_client.exceptionsimport (
ArcError, # Base exceptionArcConnectionError, # Connection failuresArcAuthenticationError,# Auth failures (401)ArcQueryError, # Query execution errorsArcIngestionError, # Write failuresArcValidationError, # Invalid inputArcNotFoundError, # Resource not found (404)ArcRateLimitError, # Rate limited (429)ArcServerError, # Server errors (5xx)
)
try:
client.query.query("INVALID SQL")
exceptArcQueryErrorase:
print(f"Query failed: {e}")
exceptArcConnectionErrorase:
print(f"Connection failed: {e}")

Documentation

See the full documentation for detailed guides and API reference.

License

MIT License - see LICENSE for details.

About

Python SDK for Arc time-series database.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages