Skip to content
Open
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
2 changes: 1 addition & 1 deletion .github/workflows/build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -329,7 +329,7 @@ jobs:
BACKPORT_TARGET_BRANCH: ${{ inputs.backport_target_branch }}
run: |
# Backport builds filter by the checked-out build.sbt.
want=(DAO Auth Config Resource Util PyBuilder WorkflowCore
want=(DAO Auth Config Observability Resource Util PyBuilder WorkflowCore
WorkflowOperator WorkflowCompiler WorkflowExecutionService)
tasks=()
if [ -n "${BACKPORT_TARGET_BRANCH}" ]; then
Expand Down
20 changes: 20 additions & 0 deletions access-control-service/LICENSE-binary
Original file line number Diff line number Diff line change
Expand Up @@ -242,6 +242,9 @@ Scala/Java jars:
- com.google.guava.listenablefuture-9999.0-empty-to-avoid-conflict-with-guava.jar
- com.google.j2objc.j2objc-annotations-2.8.jar
- com.helger.profiler-1.1.1.jar
- com.squareup.okhttp3.okhttp-4.12.0.jar
- com.squareup.okio.okio-3.6.0.jar
- com.squareup.okio.okio-jvm-3.6.0.jar
- com.thesamet.scalapb.lenses_2.13-0.11.20.jar
- com.thesamet.scalapb.scalapb-json4s_2.13-0.12.0.jar
- com.thesamet.scalapb.scalapb-runtime_2.13-0.11.20.jar
Expand Down Expand Up @@ -300,6 +303,18 @@ Scala/Java jars:
- io.fabric8.kubernetes-model-scheduling-6.12.1.jar
- io.fabric8.kubernetes-model-storageclass-6.12.1.jar
- io.fabric8.zjsonpatch-0.3.0.jar
- io.opentelemetry.opentelemetry-api-1.50.0.jar
- io.opentelemetry.opentelemetry-context-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-common-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-otlp-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-otlp-common-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-sender-okhttp-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-common-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-extension-autoconfigure-spi-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-logs-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-metrics-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-trace-1.50.0.jar
- io.r2dbc.r2dbc-spi-1.0.0.RELEASE.jar
- jakarta.inject.jakarta.inject-api-2.0.1.jar
- jakarta.validation.jakarta.validation-api-3.0.2.jar
Expand All @@ -318,6 +333,11 @@ Scala/Java jars:
- org.hibernate.validator.hibernate-validator-7.0.5.Final.jar
- org.javassist.javassist-3.30.2-GA.jar
- org.jboss.logging.jboss-logging-3.5.3.Final.jar
- org.jetbrains.annotations-13.0.jar
- org.jetbrains.kotlin.kotlin-stdlib-1.9.10.jar
- org.jetbrains.kotlin.kotlin-stdlib-common-1.9.10.jar
- org.jetbrains.kotlin.kotlin-stdlib-jdk7-1.9.10.jar
- org.jetbrains.kotlin.kotlin-stdlib-jdk8-1.9.10.jar
- org.jooq.jooq-3.19.36.jar
- org.json4s.json4s-ast_2.13-4.0.1.jar
- org.json4s.json4s-jackson-core_2.13-4.0.1.jar
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,19 @@ server:

logging:
level: ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
loggers:
# Cap noisy frameworks at WARN so TRACE/DEBUG surfaces Texera code
# (org.apache.texera) without the framework firehose.
"org.apache.pekko": WARN
"org.apache.iceberg": WARN
"org.apache.hadoop": WARN
"org.apache.kafka": WARN
"org.eclipse.jetty": WARN
"org.glassfish.jersey": WARN
"io.grpc": WARN
"io.netty": WARN
"com.zaxxer.hikari": WARN
"software.amazon.awssdk": WARN
appenders:
- type: console
threshold: ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,8 @@ class AccessControlService extends Application[AccessControlServiceConfiguration
configuration: AccessControlServiceConfiguration,
environment: Environment
): Unit = {
// Bridge this service's logs to the OTel collector under its own service.name.
org.apache.texera.observability.OtelInit.init("access-control-service")
// Serve backend at /api
environment.jersey.setUrlPattern("/api/*")

Expand Down
12 changes: 12 additions & 0 deletions amber/LICENSE-binary-java
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,18 @@ Scala/Java jars:
- org.jspecify.jspecify-1.0.0.jar
- io.opencensus.opencensus-api-0.31.1.jar
- io.opencensus.opencensus-contrib-http-util-0.31.1.jar
- io.opentelemetry.opentelemetry-api-1.50.0.jar
- io.opentelemetry.opentelemetry-context-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-common-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-otlp-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-otlp-common-1.50.0.jar
- io.opentelemetry.opentelemetry-exporter-sender-okhttp-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-common-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-extension-autoconfigure-spi-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-logs-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-metrics-1.50.0.jar
- io.opentelemetry.opentelemetry-sdk-trace-1.50.0.jar
- io.perfmark.perfmark-api-0.27.0.jar
- io.r2dbc.r2dbc-spi-1.0.0.RELEASE.jar
- io.reactivex.rxjava3.rxjava-3.1.12.jar
Expand Down
12 changes: 12 additions & 0 deletions amber/src/main/resources/computing-unit-master-config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,18 @@ logging:
level: ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
loggers:
"io.dropwizard": ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
# Cap noisy frameworks at WARN so TRACE/DEBUG surfaces Texera code
# (org.apache.texera) without the framework firehose.
"org.apache.pekko": WARN
"org.apache.iceberg": WARN
"org.apache.hadoop": WARN
"org.apache.kafka": WARN
"org.eclipse.jetty": WARN
"org.glassfish.jersey": WARN
"io.grpc": WARN
"io.netty": WARN
"com.zaxxer.hikari": WARN
"software.amazon.awssdk": WARN
appenders:
- type: console
logFormat: "[%date{ISO8601}] [%level] [%logger] [%thread] - %msg %n"
Expand Down
12 changes: 12 additions & 0 deletions amber/src/main/resources/texera-compiling-service-web-config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,18 @@ logging:
level: ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
loggers:
"io.dropwizard": ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
# Cap noisy frameworks at WARN so TRACE/DEBUG surfaces Texera code
# (org.apache.texera) without the framework firehose.
"org.apache.pekko": WARN
"org.apache.iceberg": WARN
"org.apache.hadoop": WARN
"org.apache.kafka": WARN
"org.eclipse.jetty": WARN
"org.glassfish.jersey": WARN
"io.grpc": WARN
"io.netty": WARN
"com.zaxxer.hikari": WARN
"software.amazon.awssdk": WARN
appenders:
- type: console
logFormat: "[%date{ISO8601}] [%level] [%logger] [%thread] - %msg %n"
Expand Down
12 changes: 12 additions & 0 deletions amber/src/main/resources/web-config.yml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,18 @@ logging:
level: ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
loggers:
"io.dropwizard": ${TEXERA_SERVICE_LOG_LEVEL:-INFO}
# Cap noisy frameworks at WARN so TRACE/DEBUG surfaces Texera code
# (org.apache.texera) without the framework firehose.
"org.apache.pekko": WARN
"org.apache.iceberg": WARN
"org.apache.hadoop": WARN
"org.apache.kafka": WARN
"org.eclipse.jetty": WARN
"org.glassfish.jersey": WARN
"io.grpc": WARN
"io.netty": WARN
"com.zaxxer.hikari": WARN
"software.amazon.awssdk": WARN
appenders:
- type: console
logFormat: "[%date{ISO8601}] [%level] [%logger] [%thread] - %msg %n"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,8 @@ class ComputingUnitMaster extends io.dropwizard.Application[Configuration] with
}

override def run(configuration: Configuration, environment: Environment): Unit = {
org.apache.texera.observability.OtelInit.init("computing-unit-master")
org.apache.texera.web.observability.WorkflowMetricsRecorder.init()
ObjectMapperUtils.warmupObjectMapperForOperatorsSerde()

SqlServer.initConnection(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,7 @@ class TexeraWebApplication
}

override def run(configuration: TexeraWebConfiguration, environment: Environment): Unit = {
org.apache.texera.observability.OtelInit.init("texera-web-application")
ObjectMapperUtils.warmupObjectMapperForOperatorsSerde()

// serve backend at /api
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.texera.web.observability

import com.typesafe.scalalogging.LazyLogging
import org.apache.texera.amber.core.virtualidentity.ExecutionIdentity
import org.apache.texera.amber.engine.architecture.rpc.controlreturns.WorkflowAggregatedState
import org.apache.texera.amber.engine.architecture.rpc.controlreturns.WorkflowAggregatedState.{
COMPLETED,
FAILED,
KILLED,
PAUSED,
PAUSING,
RESUMING,
RUNNING
}
import org.apache.texera.observability.WorkflowMetrics
import org.apache.texera.observability.WorkflowMetrics.WorkflowKind
import org.apache.texera.web.service.WorkflowService

import java.util.concurrent.ConcurrentHashMap

/**
* Drives [[WorkflowMetrics]] from amber's execution lifecycle. The metric
* instruments live in common/config; this object is the single place
* that records them, so the lifecycle code only needs one-line calls.
*
* - start / terminal counters and the duration histogram are recorded
* from [[onStart]] and [[onStateChange]];
* - the always-polled `texera.workflow.active` gauge is sourced from the
* live WorkflowService registry via the supplier registered in [[init]].
*/
object WorkflowMetricsRecorder extends LazyLogging {

// Start time + kind per in-flight run, so a terminal transition can emit
// a duration and attribute the outcome. Keyed by the execution identity.
private val inFlight = new ConcurrentHashMap[ExecutionIdentity, (Long, WorkflowKind)]()

private val ActiveStates: Set[WorkflowAggregatedState] = Set(RUNNING, PAUSING, PAUSED, RESUMING)
private val TerminalStates: Set[WorkflowAggregatedState] = Set(COMPLETED, FAILED, KILLED)

/** Register the active-executions gauge supplier and bind the instruments.
* Call once at startup, after OtelInit.init. The supplier is polled on
* every metric collection, so it must never throw.
*/
def init(): Unit = {
WorkflowMetrics.setActiveExecutionsSupplier(() =>
try {
WorkflowService.getAllWorkflowServices.iterator
.flatMap(s => Option(s.executionService.getValue))
.map(_.executionStateStore.metadataStore.getState.state)
.count(ActiveStates.contains)
.toLong
} catch {
case _: Throwable => 0L
}
)
WorkflowMetrics.ensureBound()
}

/** Record that a run started. */
def onStart(
executionId: ExecutionIdentity,
kind: WorkflowKind = WorkflowKind.Interactive
): Unit = {
inFlight.put(executionId, (System.currentTimeMillis(), kind))
WorkflowMetrics.recordStart(kind)
}

/** Record terminal counters + duration exactly once, on the first
* transition from a non-terminal into a terminal state. Safe to call on
* every state change; non-terminal and repeat-terminal calls are no-ops.
*/
def onStateChange(
executionId: ExecutionIdentity,
oldState: WorkflowAggregatedState,
newState: WorkflowAggregatedState
): Unit = {
if (!TerminalStates.contains(newState) || TerminalStates.contains(oldState)) return
val entry = Option(inFlight.remove(executionId))
val kind = entry.map(_._2).getOrElse(WorkflowKind.Interactive)
val durationSec = entry.map(e => (System.currentTimeMillis() - e._1) / 1000.0).getOrElse(0.0)
newState match {
case COMPLETED => WorkflowMetrics.recordCompletion(kind, durationSec)
case FAILED => WorkflowMetrics.recordFailure(kind, durationSec)
case KILLED => WorkflowMetrics.recordCancellation(kind)
case _ => ()
}
}
}
Loading
Loading