Skip to content

Logging

DataCoolie has two independent streams:

  • SystemLogger captures framework Python records for operational diagnosis.
  • ExecutionLogger records terminal dataflow observations and one mutable job summary.

Both streams use the same bounded JSON Lines persistence writer. The default is a snapshot: one remote .json file per stream is replaced with the latest complete snapshot. persistence_mode="batch" publishes immutable *_part_00000001.json files after the configured size or time trigger. Files use the .json extension for platform preview compatibility, but contain one compact JSON object per line.

SystemLogger

SystemLogger is inert until activate(). Driver activates it (and then ExecutionLogger) after pure configuration/provider binding but before provider metadata I/O, so provider startup records are captured. log_level controls console output; file_level controls captured records. Capture diagnostics remain console-only so a storage failure cannot recursively write to the same sink.

Each system record includes log_schema_version, then _type="system_log", then datacoolie_version (the installed producer package version), followed by log_session_id, job_id, job_num, and job_index, then the LogRecord projection (ts, level, logger, msg, and optional source location, exception, event_name, dataflow_id, and dataflow_run_id). Empty streams do not create a file.

Framework lifecycle anchors use stable dotted event names, including session.starting, session.ready, session.startup_failed, session.finishing, operation.started, operation.finished, dataflow.started, and dataflow.finished. The two dataflow identifiers are bound only for the active execution scope and are restored when work returns; rejected operations do not emit a started event. Event production never forces a remote flush.

The persisted record keeps dataflow_id, dataflow_run_id, and event_name as independent fields. Console formatting adds them only as an optional display suffix: dataflow_id is shown in its own brackets and dataflow_run_id:event_name in a second bracket, omitting any value that is absent. For example, a fully correlated record is rendered as [orders] [run-123:dataflow.started]. This is presentation-only; the original message remains the semantic message and is not rewritten with identifiers, so structured consumers and persistence can query each field independently.

ExecutionLogger

Job and dataflow records also include the automatic datacoolie_version field. All three record kinds carry log_session_id. Driver assigns the same value to both loggers for its session, so reuse of a caller-supplied job_id does not erase session identity. Correlate individual executions with dataflow_run_id; replay chunks have their own execution IDs. Package releases and log schema versions evolve independently: readers ignore unknown fields, and compatible optional additions keep the same schema version. Readers of historical logs must tolerate an absent producer version.

Execution logging accepts only terminal succeeded, failed, or skipped runtime observations. Preparation failures are terminal failed rows. An execution killed before it is admitted to a dataflow/chunk boundary has no fabricated row; if the scheduler itself fails after admitting a dataflow, the boundary creates one failed runtime for that dataflow so the job accounting remains truthful. File persistence requires an activated logger with a platform and output path. activate() supplies a default DataCoolieRunConfig when a standalone logger has not been given one; Driver-owned loggers receive the Driver's validated configuration before activation. Uploads and final close are bounded best-effort operations, so a job summary may be absent or stale after a storage failure. See Logging layout. The dataflow record keeps declarative metadata, including the original source_query. The flattened source_action JSON string may contain the exact SQL sent to the engine, including runtime predicates.

The job record uses _type="job_run_log"; dataflow records use _type="dataflow_run_log". The job record is always a replace-one snapshot, even when dataflow persistence uses batch mode. It contains aggregate counters, component names, lifecycle status, log_records_dropped, log_bytes_dropped, message_truncated, and the optional caller-owned run_attributes string. Call ExecutionLogger.finish_job(...) to provide the business outcome when using the logger without a Driver; close() alone does not invent success.

Every terminal dataflow and phase runtime uses one optional message field. For failed terminal observations, the JobRuntime message is a compact index of failed dataflow identities (name [id], with an ID-only fallback), joined with ;. For skipped observations, the same field carries the human-readable skip reason. The dataflow row retains the full phase details. Driver-owned session failures (for example startup, scheduler-contract, or provider-teardown failures) are appended only when they are independent of an already observed terminal dataflow, so one failure is not counted twice across the two summary owners. System records retain the event context and complete exception chain; a concise message plus the traceback's exception line is intentional when both aid diagnosis.

In flattened dataflow records, phase explanations are exposed as source_message, transform_message, and destination_message. The top-level message remains the explanation for the dataflow outcome; the phase fields retain more specific partial-failure evidence.

Layout

When created by Driver, log_base_path is split into the two component roots below. For standalone loggers, LogConfig.output_path is already that logger's component root.

execution_logs/
├── job_run_log/<partition>/job_<stem>.json
└── dataflow_run_log/<partition>/dataflow_<stem>.json
system_logs/<partition>/system_<stem>.json

Batch system/dataflow streams add _part_<sequence:08d> before .json. The stem contains the job-start timestamp, job number/index, and a filesystem-safe job-id token. The raw job id remains in each record alongside the session log_session_id. Date partitioning follows LogConfig.partition_pattern and batch parts use their sealing time. Snapshot paths are frozen at activation.

Run attributes

DataCoolieRunConfig.run_attributes is an optional JSON object supplied by the caller for correlation with an external orchestrator, for example:

DataCoolieRunConfig(
    job_id="orders",
    run_attributes={"data_factory_run_id": "adf-123", "glue_job_run_id": "jr-456"},
)

Keys are caller-defined and are not interpreted by the framework. Objects, arrays, finite numbers, booleans, strings, and null are accepted; raw JSON strings, non-string keys, cycles, and non-finite numbers fail configuration validation. Attributes are serialized once, deterministically, in the job summary and are not repeated in every dataflow/system record.

Configuration

Applications configure logging through the two public logger types. The system logger owns operational capture and the execution logger owns structured run records:

from datacoolie.core.models.run_config import DataCoolieRunConfig
from datacoolie.logging import ExecutionLogger, LogConfig, SystemLogger
from datacoolie.platforms.local_platform import LocalPlatform

platform = LocalPlatform()
run_config = DataCoolieRunConfig(
    job_id="standalone-logging",
    run_attributes={"owner": "example"},
)
log_config = LogConfig(output_path="logs")
system_logger = SystemLogger(log_config, platform)
execution_logger = ExecutionLogger(log_config, platform)
system_logger.set_run_config(run_config)
execution_logger.set_run_config(run_config)

# Nested contexts activate system capture first and close execution first.
with system_logger:
    with execution_logger:
        # Run the standalone work here and emit terminal dataflow rows as needed.
        execution_logger.finish_job("succeeded")

SystemLogger.activate() claims the process-wide capture session. Driver activates an injected or auto-created SystemLogger before provider initialization, then activates the execution logger and emits the session-starting anchor. An ExecutionLogger used without a SystemLogger does not implicitly configure console/capture output. Invalid levels or modes fail before handlers or an active capture owner are changed.

Important LogConfig fields:

  • persistence_mode: "snapshot" (default) or "batch" for system/dataflow streams.
  • flush_interval_seconds: 300 by default; 0 disables time-triggered flushes (batch size wakeups and final close still work).
  • flush_batch_bytes: approximately 4 MiB by default in batch mode.
  • buffer_memory_bytes / spool_max_bytes: bounded encoded-data budgets (64 MiB / 512 MiB defaults). The writer reserves space for both retained and temporary upload copies, so the practical retained limit is lower than spool_max_bytes.
  • spool_directory: optional local spool directory.
  • close_timeout_seconds: bounded terminal sink wait (10 seconds by default).
  • storage_mode: controls only the internal capture buffer ("memory" or "file").

At capacity, new records are dropped and counted in persistence statistics; business execution does not fail. A failed upload keeps the exact frozen payload and destination for a later retry; records admitted after that failure form a later batch. Startup and close waits are bounded, and a timed-out or skipped write is never reported as successful. Terminal close drains already-admitted parts while its shared deadline allows; a sink that remains in flight is reported as timed out. The platform owns cloud replacement semantics, so cross-file atomicity and exactly-once delivery are not claimed.

The logging contract is framework-first. DataCoolie Studio consumers migrate after framework fixtures and schema/layout verification are complete.