Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 13 additions & 18 deletions api/continuous_query_routes.py
Original file line number Diff line number Diff line change
Expand Up @@ -669,34 +669,29 @@ async def execute_continuous_query(
for row in data:
record = dict(zip(columns, row))

# Ensure time is in milliseconds (arrow_writer expects ms, not microseconds)
# Convert time to datetime (arrow_writer will store as timestamp[us])
# OPTIMIZATION: When continuous queries use epoch_us() in SQL, time arrives as integer microseconds
time_milliseconds = None
if 'time' in record:
time_val = record['time']
if isinstance(time_val, (int, float)):
# Integer time - determine format by magnitude
if time_val < 1_000_000_000_000: # Less than 1T = seconds or milliseconds
if time_val < 10_000_000_000: # Less than 10B = seconds (e.g., 1729780800)
time_milliseconds = int(time_val * 1_000)
else: # Between 10B and 1T = milliseconds
time_milliseconds = int(time_val)
else:
# Greater than 1T = microseconds (e.g., 1729780800000000)
time_milliseconds = int(time_val / 1_000)
elif isinstance(time_val, datetime):
time_milliseconds = int(time_val.timestamp() * 1_000)
# Integer time - determine format by magnitude and convert to datetime
if time_val < 1e10: # Less than 10B = seconds (e.g., 1729780800)
record['time'] = datetime.fromtimestamp(time_val, tz=timezone.utc)
elif time_val < 1e13: # Between 10B and 1T = milliseconds
record['time'] = datetime.fromtimestamp(time_val / 1_000, tz=timezone.utc)
else: # Greater than 1T = microseconds (e.g., 1729780800000000)
record['time'] = datetime.fromtimestamp(time_val / 1_000_000, tz=timezone.utc)
elif isinstance(time_val, str):
# LEGACY: Parse timestamp string from DuckDB (less efficient)
# Recommend using epoch_us() in queries instead
from dateutil import parser
dt = parser.parse(time_val)
time_milliseconds = int(dt.timestamp() * 1_000)

record['time'] = time_milliseconds
record['time'] = parser.parse(time_val)
# elif isinstance(time_val, datetime): already correct, keep as-is

# Create flat record (arrow_writer expects flat dictionaries, not nested tags/fields)
final_time = record.get('time', int(datetime.utcnow().timestamp() * 1_000))
final_time = record.get('time')
if not isinstance(final_time, datetime):
final_time = datetime.utcnow().replace(tzinfo=timezone.utc)

flat_record = {
'measurement': query['destination_measurement'],
Expand Down
28 changes: 24 additions & 4 deletions ingest/arrow_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,10 +160,16 @@ def _infer_schema(self, columns: Dict[str, List]) -> pa.Schema:
Infer Arrow schema from column data.

Handles:
- Timestamps (datetime → timestamp[ms])
- Timestamps (datetime → timestamp[us] - microsecond precision)
- Integers (int64)
- Floats (float64)
- Strings (utf8)

Note: Microsecond precision is the industry standard for observability:
- Sufficient for distributed tracing (sub-millisecond spans)
- Compatible with most databases and tools
- No storage overhead vs milliseconds (both 64-bit)
- DuckDB native format (no conversion needed)
"""
fields = []

Expand All @@ -175,8 +181,9 @@ def _infer_schema(self, columns: Dict[str, List]) -> pa.Schema:
# All None, default to string
arrow_type = pa.string()
elif isinstance(sample, datetime):
# Timestamp with millisecond precision
arrow_type = pa.timestamp('ms')
# Timestamp with microsecond precision (us = 10^-6 seconds)
# Industry standard for observability and time series databases
arrow_type = pa.timestamp('us')
elif isinstance(sample, bool):
arrow_type = pa.bool_()
elif isinstance(sample, int):
Expand Down Expand Up @@ -509,7 +516,20 @@ async def _flush_records(self, measurement: str, records: List[Dict[str, Any]]):
if isinstance(record['time'], str):
record['time'] = datetime.fromisoformat(record['time'].replace('Z', '+00:00'))
elif isinstance(record['time'], (int, float)):
record['time'] = datetime.fromtimestamp(record['time'] / 1000) # Assume ms
# Auto-detect timestamp unit based on magnitude:
# - Less than 1e10: seconds (e.g., 1730246400)
# - Between 1e10 and 1e13: milliseconds (e.g., 1730246400000)
# - Greater than 1e13: microseconds (e.g., 1730246400000000)
timestamp_val = record['time']
if timestamp_val < 1e10:
# Seconds
record['time'] = datetime.fromtimestamp(timestamp_val)
elif timestamp_val < 1e13:
# Milliseconds
record['time'] = datetime.fromtimestamp(timestamp_val / 1000)
else:
# Microseconds
record['time'] = datetime.fromtimestamp(timestamp_val / 1000000)

# Get time range for filename
times = [r['time'] for r in records if 'time' in r]
Expand Down
41 changes: 29 additions & 12 deletions ingest/msgpack_decoder.py
Original file line number Diff line number Diff line change
Expand Up @@ -147,10 +147,19 @@ def _decode_columnar(self, obj: Dict[str, Any]) -> Dict[str, Any]:
if 'time' in columns:
time_col = columns['time']
if time_col and isinstance(time_col[0], (int, float)):
columns['time'] = [
datetime.fromtimestamp(t / 1000, tz=timezone.utc)
for t in time_col
]
# Auto-detect timestamp unit based on magnitude (same as arrow_writer)
def convert_timestamp(t):
if t < 1e10:
# Seconds
return datetime.fromtimestamp(t, tz=timezone.utc)
elif t < 1e13:
# Milliseconds (default for msgpack API)
return datetime.fromtimestamp(t / 1000, tz=timezone.utc)
else:
# Microseconds
return datetime.fromtimestamp(t / 1000000, tz=timezone.utc)

columns['time'] = [convert_timestamp(t) for t in time_col]

# Return columnar record marker
return {
Expand All @@ -176,14 +185,22 @@ def _decode_single(self, obj: Dict[str, Any]) -> Dict[str, Any]:
if isinstance(measurement, int):
measurement = f"measurement_{measurement}" # TODO: lookup from registry

# Extract timestamp (milliseconds)
timestamp_ms = obj.get('t')
if timestamp_ms is None:
# Use current time if not provided
timestamp_ms = int(datetime.now(timezone.utc).timestamp() * 1000)

# Convert to datetime
timestamp = datetime.fromtimestamp(timestamp_ms / 1000, tz=timezone.utc)
# Extract timestamp (default: milliseconds for backwards compatibility)
timestamp_val = obj.get('t')
if timestamp_val is None:
# Use current time if not provided (generate as milliseconds for compatibility)
timestamp_val = int(datetime.now(timezone.utc).timestamp() * 1000)

# Convert to datetime with auto-detection
if timestamp_val < 1e10:
# Seconds
timestamp = datetime.fromtimestamp(timestamp_val, tz=timezone.utc)
elif timestamp_val < 1e13:
# Milliseconds (default for msgpack API)
timestamp = datetime.fromtimestamp(timestamp_val / 1000, tz=timezone.utc)
else:
# Microseconds
timestamp = datetime.fromtimestamp(timestamp_val / 1000000, tz=timezone.utc)

# Extract host
host = obj.get('h', 'unknown')
Expand Down