Data Flow Application Logging
Efficient log processing is essential for debugging, monitoring, and performance tuning of distributed Apache Spark applications, especially in cloud native environments such as Data Flow.
A Spark job emits logs from several components, primarily the driver and executors, each generating separate stdout and stderr
streams.
-
The driver logs capture initialization, Spark context setup, job planning, and application-level messages (for example,
log.infofrom user code). These logs are critical for understanding the job lifecycle and contain high-level orchestration events, application logs, and configuration metadata. -
The executor logs, on the other hand, reflect distributed task execution, including per-task
stdoutoutputs and error diagnostics (for example, stack traces and memory errors). Each executor independently writes to its ownstdoutandstderrstreams, which might produce thousands of log fragments per job, especially under parallelism.
In addition to Spark-internal logs, application-level stdout and stderr (for example, those produced using SLF4J,
println, or System.out or System.err in user-defined transformations) must be captured and correlated with
Spark stages and tasks. Depending on the OCI Object Storage, these logs can reside in different storage locations.
Effective log processing requires:
-
Aggregating and timestamp-aligning logs from several sources.
-
Tagging logs with identifiers (stage, task, executor ID).
-
Linking
stdoutorstderracross retries and executor failures.
Data Flow Implementation
Data Flow provides automated mechanisms to persist driver and executor logs into OCI Object Storage, with hierarchical organization and retention policies. These logs are accessible through the OCI Console or CLI.
Besides this standard method, Data Flow supports integration with OCI Logging to capture both Spark diagnostic and custom application logs. To enable this functionality, users must ensure the appropriate Identity and Access Management (IAM) policies are in place. This includes:
-
Applying policies that grant access to OCI Logging services. For more information, see Oracle Cloud Infrastructure Logging Policies. for detailed instructions.
-
Setting up IAM policies specific to Data Flow, as outlined in Set Up Identity and Access Management Policies., includes permissions required for enhanced logging.
Without these policies, logging capabilities might be limited or unavailable. Proper configuration ensures complete visibility into Spark job execution and application behavior.
- Spark Diagnostic Logs (Service Logs)
-
Spark Diagnostic Logs are emitted by the underlying Spark infrastructure managed by Data Flow. These logs include output from the Spark Driver and Executors, such as:
-
Task execution events
-
Resource usage and memory management
-
Error and exception traces
When enabled, all runs started from the associated Data Flow application automatically generate diagnostic logs for the Spark Driver and Executors. These logs are classified as Service Logs because they originate from an Oracle Cloud Infrastructure native service, in this case, Data Flow.
-
- Application Logs (Custom Logs)
-
Application Logs capture messages emitted directly from the Spark application logic, typically through logging frameworks such as
SLF4J, or standard output. These logs might include:-
Business process checkpoints
-
Custom debug or info messages
-
Processing metrics and counters
-
Exceptions and user-defined error handling
Unlike diagnostic logs, application logs are written from within the Spark code and give insight into domain-specific logic and execution flow.
-
After the OCI Logging is configured in the application, in a standard execution run, Data Flow generates four sets of OCI Logging entries. The Archive Logs are generated after the fact from the OCI Logging entries and take a few minutes to be delivered to the assigned bucket in the customer tenancy.
Data Flow Application Logs
In Spark applications where business process observability is required, for example, tracking custom milestones, domain events, or application-specific metrics, Data Flow exposes application-level log output through internal mechanisms, ensuring standard output is surfaced in a consistent format.
Establishing the communication channel within the Spark context is essential for reliably capturing and emitting business process logic within a Spark application. This ensures that the channel is initialized alongside Spark's runtime and is accessible to components such as listeners, drivers, or executors during job execution. Spark jobs are distributed among the workers and the driver and are event-driven. Lifecycle events occur only after the Spark context is active. Therefore, within Data Flow implementations, the instrumented business logic needs the Spark context to be fully initialized to receive and respond to these events.
Considerations of High-Traffic Regions
In high-traffic OCI regions, diagnostic and application logging, especially to OCI Object Storage, can lead to throttled APIs and delayed job observability.
Where possible, prefer OCI Logging over raw Object Storage for diagnostic and application logs. The Logging service is designed for real-time, high-throughput log ingestion.
Regarding the development lifecycle, a Spark or general application can maintain different logging levels across development, test, and production
environments by combining configuration-based logging frameworks such as SLF4J with environment-aware code practices. Externalizing
log-level settings into a logging configuration file and selecting the appropriate configuration based on the environment can create a development lifecycle
environment (for example, development, test, and production) that respects the region's current processing capabilities.
Most open source libraries (including Spark) use SLF4J for logging. SLF4J lets the application remain agnostic to the
underlying logging system, which is great for libraries or applications running in different environments. Using SLF4J ensures consistent
integration and fewer conflicts across Data Flow dependencies. SLF4J natively supports parameter
substitution, which is more efficient than string interpolation or concatenation (for example, no unnecessary string creation when DEBUG is disabled).
A new logging class, such as Console Logger, might be relevant to implement within these characteristics, as in the body of the following code:
protected lazy val LOG: org.slf4j.Logger =
ConsoleLogger(getClass.getName)//Copyright (c) 2025 Oracle and/or its affiliates.
//The Universal Permissive License (UPL), Version 1.0
package com.oracle.delta
import org.slf4j.{Logger, LoggerFactory, Marker}
class ConsoleLogger(private val underlying: Logger) extends Logger {
override def getName: String = underlying.getName
override def isTraceEnabled: Boolean = underlying.isTraceEnabled
override def isTraceEnabled(marker: Marker): Boolean = underlying.isTraceEnabled(marker)
// --- TRACE ---
override def trace(marker: Marker, msg: String): Unit = trace(msg)
override def trace(marker: Marker, format: String, arg: Any): Unit = trace(format, arg)
override def trace(format: String, arg: Any): Unit =
trace(formatMsg(format, Seq(arg)))
override def trace(msg: String): Unit = {
logConsole("TRACE", msg)
underlying.trace(msg)
}
override def trace(marker: Marker, format: String, arg1: Any, arg2: Any): Unit = trace(format, arg1, arg2)
override def trace(format: String, arg1: Any, arg2: Any): Unit =
trace(formatMsg(format, Seq(arg1, arg2)))
override def trace(marker: Marker, format: String, argArray: AnyRef*): Unit = trace(format, argArray: _*)
override def trace(format: String, arguments: AnyRef*): Unit =
trace(formatMsg(format, arguments))
override def trace(marker: Marker, msg: String, t: Throwable): Unit = trace(msg, t)
override def trace(msg: String, t: Throwable): Unit = {
logConsole("TRACE", s"$msg - ${t.getMessage}")
underlying.trace(msg, t)
}
override def isDebugEnabled: Boolean = underlying.isDebugEnabled
override def isDebugEnabled(marker: Marker): Boolean = underlying.isDebugEnabled(marker)
// --- DEBUG ---
override def debug(marker: Marker, msg: String): Unit = debug(msg)
override def debug(msg: String): Unit = {
logConsole("DEBUG", msg)
underlying.debug(msg)
}
override def debug(marker: Marker, format: String, arg: Any): Unit = debug(format, arg)
override def debug(format: String, arg: Any): Unit =
debug(formatMsg(format, Seq(arg)))
override def debug(marker: Marker, format: String, arg1: Any, arg2: Any): Unit = debug(format, arg1, arg2)
override def debug(format: String, arg1: Any, arg2: Any): Unit =
debug(formatMsg(format, Seq(arg1, arg2)))
override def debug(marker: Marker, format: String, arguments: AnyRef*): Unit = debug(format, arguments: _*)
override def debug(format: String, arguments: AnyRef*): Unit =
debug(formatMsg(format, arguments))
override def debug(marker: Marker, msg: String, t: Throwable): Unit = debug(msg, t)
override def debug(msg: String, t: Throwable): Unit = {
logConsole("DEBUG", s"$msg - ${t.getMessage}")
underlying.debug(msg, t)
}
override def isInfoEnabled: Boolean = underlying.isInfoEnabled
override def isInfoEnabled(marker: Marker): Boolean = underlying.isInfoEnabled(marker)
// --- INFO ---
override def info(marker: Marker, msg: String): Unit = info(msg)
override def info(marker: Marker, format: String, arg: Any): Unit = info(format, arg)
override def info(format: String, arg: Any): Unit =
info(formatMsg(format, Seq(arg)))
private def formatMsg(format: String, args: Seq[Any]): String =
args.foldLeft(format)((msg, arg) => msg.replaceFirst("\\{\\}", arg.toString))
override def info(msg: String): Unit = {
logConsole("INFO", msg)
underlying.info(msg)
}
private def logConsole(level: String, msg: String): Unit = println(s"[$level] $msg")
override def info(marker: Marker, format: String, arg1: Any, arg2: Any): Unit = info(format, arg1, arg2)
override def info(format: String, arg1: Any, arg2: Any): Unit =
info(formatMsg(format, Seq(arg1, arg2)))
override def info(marker: Marker, format: String, arguments: AnyRef*): Unit = info(format, arguments: _*)
override def info(format: String, arguments: AnyRef*): Unit =
info(formatMsg(format, arguments))
override def info(marker: Marker, msg: String, t: Throwable): Unit = info(msg, t)
override def info(msg: String, t: Throwable): Unit = {
logConsole("INFO", s"$msg - ${t.getMessage}")
underlying.info(msg, t)
}
// --- WARN ---
override def isWarnEnabled: Boolean = underlying.isWarnEnabled
override def isWarnEnabled(marker: Marker): Boolean = underlying.isWarnEnabled(marker)
override def warn(marker: Marker, msg: String): Unit = warn(msg)
override def warn(marker: Marker, format: String, arg: Any): Unit = warn(format, arg)
override def warn(format: String, arg: Any): Unit =
warn(formatMsg(format, Seq(arg)))
override def warn(marker: Marker, format: String, arg1: Any, arg2: Any): Unit = warn(format, arg1, arg2)
override def warn(format: String, arg1: Any, arg2: Any): Unit =
warn(formatMsg(format, Seq(arg1, arg2)))
override def warn(msg: String): Unit = {
logConsole("WARN", msg)
underlying.warn(msg)
}
override def warn(marker: Marker, format: String, arguments: AnyRef*): Unit = warn(format, arguments: _*)
override def warn(format: String, arguments: AnyRef*): Unit =
warn(formatMsg(format, arguments))
override def warn(marker: Marker, msg: String, t: Throwable): Unit = warn(msg, t)
override def warn(msg: String, t: Throwable): Unit = {
logConsole("WARN", s"$msg - ${t.getMessage}")
underlying.warn(msg, t)
}
// --- ERROR ---
override def isErrorEnabled: Boolean = underlying.isErrorEnabled
override def isErrorEnabled(marker: Marker): Boolean = underlying.isErrorEnabled(marker)
override def error(marker: Marker, msg: String): Unit = error(msg)
override def error(marker: Marker, format: String, arg: Any): Unit = error(format, arg)
override def error(format: String, arg: Any): Unit =
error(formatMsg(format, Seq(arg)))
override def error(marker: Marker, format: String, arg1: Any, arg2: Any): Unit = error(format, arg1, arg2)
override def error(format: String, arg1: Any, arg2: Any): Unit =
error(formatMsg(format, Seq(arg1, arg2)))
override def error(marker: Marker, format: String, arguments: AnyRef*): Unit = error(format, arguments: _*)
override def error(format: String, arguments: AnyRef*): Unit =
error(formatMsg(format, arguments))
override def error(msg: String): Unit = {
logConsole("ERROR", msg)
underlying.error(msg)
}
override def error(marker: Marker, msg: String, t: Throwable): Unit = error(msg, t)
override def error(msg: String, t: Throwable): Unit = {
logConsole("ERROR", s"$msg - ${t.getMessage}")
underlying.error(msg, t)
}
}
object ConsoleLogger {
def apply(cls: Class[_]): ConsoleLogger =
new ConsoleLogger(LoggerFactory.getLogger(cls.getName))
def apply(name: String): ConsoleLogger =
new ConsoleLogger(LoggerFactory.getLogger(name))
}Recommendations
Data Flow supports two main types of logs using the OCI Logging service:
-
Spark Diagnostic Logs (service logs): Captured from the Spark driver and executor; enabled at the Data Flow application level only and apply to all later runs. These logs offer visibility into Spark runtime events such as stage execution, memory usage, and failures.
-
Application Logs: Custom logs emitted from user code (for example, using
SLF4J) that trace business logic, custom exceptions, and domain workflows.
Data Flow exposes application-level log output through internal mechanisms, ensuring standard output is surfaced in a consistent format. This is useful for finding the outcome of intermediate ETL processing stages by printing summaries of the processing or verifying application processing checkpoints.
For high-traffic regions, if further specialized logging is necessary, implementing a Console Logger can help separate the Data Flow service entries from the Application, especially for the entries associated with the cluster environment and processing necessary to troubleshoot errors and performance.