How Finelog's LogClient Writes Batches and Queries Structured Logs in Marin

Finelog's LogClient provides a high-level Python façade that converts protobuf LogEntry messages into Arrow batches for asynchronous writing, while offering both SQL-based analytics and dedicated RPC fetching for structured log retrieval.

Finelog's LogClient serves as the primary interface for Marin applications such as Iris workers and Zephyr pipelines to record and retrieve structured logs. Located in lib/finelog/src/finelog/client/log_client.py, this client manages a sophisticated two-stage architecture: asynchronously persisting log batches to a privileged log namespace and executing dual-path queries through lazily cached RPC stubs. Understanding the internal mechanics of batch enqueueing, background flushing, and query routing is essential for building reliable observability pipelines in the Marin ecosystem.

Writing Batches to the Log Namespace

The write path transforms protobuf messages into lightweight rows, buffers them in memory-backed queues, and flushes them as Arrow RecordBatches via a background daemon thread.

Connecting and Initializing the Client

Applications instantiate LogClient through the connect factory method, which configures the server endpoint and optional interceptors. This method builds the client configuration without immediately establishing RPC channels.

from finelog.client import LogClient

client = LogClient.connect("http://finelog:8080")

According to the Marin source code, LogClient.connect handles endpoint resolution and client initialization at lines 525-563.

Enqueuing Log Entries with write_batch

The write_batch method accepts a logical key and a list of protobuf LogEntry messages. It validates the input, converts messages to row objects, and delegates to an internal Table instance representing the privileged log namespace.

from finelog.rpc import logging_pb2

entries = [
    logging_pb2.LogEntry(source="my-job", data=b"event-1", level=logging_pb2.INFO),
    logging_pb2.LogEntry(source="my-job", data=b"event-2", level=logging_pb2.INFO),
]

client.write_batch(key="job-1234", messages=entries)

Internally, write_batch invokes _log_entries_to_rows to construct simple-namespace objects containing key, source, data, epoch_ms, and level fields (lines 554-566). The rows are then passed to Table.write (lines 525-553).

Background Flushing and Retry Logic

Each Table instance maintains an internal queue and a daemon thread (_run) that periodically flushes accumulated rows. The queue triggers a flush when it reaches DEFAULT_BATCH_ROWS (10,000 entries) or a 16 MiB byte threshold.

The background thread (lines 390-447) performs the following sequence:

  1. Pulls queued _PendingItem objects from the buffer
  2. Ensures namespace registration via _ensure_registered
  3. Constructs an Arrow RecordBatch via _rows_to_record_batch
  4. Transmits the batch via _send (lines 473-502)

If an RPC fails with a retryable error, the implementation uses ExponentialBackoff to re-buffer the batch and retry. Non-retryable errors result in the batch being dropped, with the drop count exposed through FlushResult.DROPPED. The final RPC transmission occurs in _stats_flush at lines 882-894.

Querying Structured Logs

LogClient exposes two distinct read paths optimized for different access patterns: fetch_logs for tail-oriented log retrieval and query for general SQL analytics.

Fetching Logs via fetch_logs (Log-Read RPC)

The fetch_logs method provides efficient, protobuf-based access to the privileged log namespace. It utilizes a purpose-built log-read RPC through LogServiceClientSync.

from finelog.rpc import logging_pb2

request = logging_pb2.FetchLogsRequest(
    namespace="log",
    start_seq=0,
    max_rows=100,
)

response = client.fetch_logs(request)
for entry in response.entries:
    print(entry.timestamp, entry.source, entry.data.decode())

The implementation at lines 613-630 lazily obtains the RPC client via _get_log_service_client and translates connection errors to domain-specific exceptions such as NamespaceNotFoundError.

Analyzing Data via query (SQL Path)

The query method enables arbitrary SQL execution against any registered namespace, including user-defined tables. It returns results as a pyarrow.Table via the stats service.

sql = '''
SELECT namespace, row_count, byte_size
FROM "finelog.stats"
WHERE namespace LIKE 'iris.%'
ORDER BY row_count DESC
LIMIT 10
'''
table = client.query(sql, max_rows=5_000)
df = table.to_pandas()

The query implementation (lines 631-648) validates row limits before invoking _stats_query (lines 665-672), which decodes the Arrow IPC stream returned by StatsServiceClientSync.query.

Client Lifecycle and Connection Management

LogClient implements sophisticated connection handling including lazy resolution, caching, and graceful shutdown to ensure reliable operation in long-running Marin pipelines.

Lazy RPC Resolution and Caching

To minimize connection overhead, _get_or_resolve_client resolves endpoints only once and caches the resulting RPC stub. Subsequent operations reuse the cached client until explicitly invalidated. This logic is implemented at lines 888-915.

Error Handling and Invalidation

When transient failures occur (such as network partitions), the client invokes _invalidate (lines 851-861) to clear the cached RPC stub. The next operation triggers a fresh endpoint resolution, allowing the client to recover from temporary service disruptions without process restarts.

Graceful Shutdown with close()

The close method ensures no pending log rows are lost during application termination. It synchronously drains every Table queue (including the internal log table) and waits for background flush threads to complete before closing underlying RPC channels.

result = client.flush()  # Optional: synchronous flush

print("Flush result:", result)  # FlushResult.SUCCEEDED, TIMEOUT, or DROPPED

client.close()

The shutdown sequence is defined at lines 582-603.

Summary

  • LogClient in lib/finelog/src/finelog/client/log_client.py serves as the primary Python interface for Marin applications to write and query structured logs.
  • write_batch converts protobuf LogEntry messages to rows and enqueues them on an internal Table, which flushes Arrow batches via a background thread when reaching 10,000 rows or 16 MiB.
  • fetch_logs provides protobuf-based RPC access optimized for tail-reading the privileged log namespace, while query enables SQL analytics across any registered namespace using Arrow IPC.
  • The client employs lazy RPC resolution with caching and automatic invalidation to handle transient failures, plus exponential backoff for retryable write errors.
  • Graceful shutdown via close() drains all pending batches, ensuring durability of in-flight log entries.

Frequently Asked Questions

How does LogClient handle retries for failed batch writes?

LogClient implements an exponential backoff strategy within the background flush thread. If _stats_flush encounters a retryable error (determined by is_retryable_error), the batch is re-buffered and the thread sleeps for an increasing duration before retrying. Non-retryable errors cause immediate batch dropping, with the failure recorded in FlushResult.DROPPED.

What is the difference between fetch_logs and query in LogClient?

fetch_logs utilizes a dedicated LogServiceClientSync RPC optimized for reading protobuf LogEntry messages from the privileged log namespace, ideal for tail-oriented log streaming. query invokes the StatsServiceClientSync to execute arbitrary SQL against any registered namespace and returns results as an Arrow IPC stream suitable for analytical workloads.

How does LogClient manage memory when buffering log batches?

The Table class buffers rows in memory and enforces two flush triggers: a row count limit of DEFAULT_BATCH_ROWS (10,000 entries) and a byte size cap of 16 MiB. When either threshold is exceeded, the background thread immediately converts the queued rows to an Arrow RecordBatch and transmits it to the server, preventing unbounded memory growth.

Can LogClient query namespaces other than the default log namespace?

Yes. While fetch_logs is restricted to the privileged log namespace, the query method can execute SQL against any registered namespace accessible through the stats service, including user-defined tables created by Iris workers or Zephyr pipelines. The method accepts standard SQL with max_rows limits to prevent excessive memory consumption.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →