Skip to content
31 changes: 23 additions & 8 deletions historyserver/cmd/collector/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"encoding/json"
"errors"
"flag"
"fmt"
"os"
Expand Down Expand Up @@ -143,6 +144,8 @@ func main() {
}
}

// RAY_COLLECTOR_ADDITIONAL_ENDPOINTS is optional: the head collector always
// polls its built-in endpoints, and anything listed here is polled on top.
var additionalEndpoints []string
if epStr := os.Getenv("RAY_COLLECTOR_ADDITIONAL_ENDPOINTS"); epStr != "" {
for _, ep := range strings.Split(epStr, ",") {
Expand All @@ -153,16 +156,20 @@ func main() {
}
}

// An unusable poll interval falls back to the default instead of exiting: the
// collector is a sidecar in the Ray head pod, so crash-looping on a bad
// observability knob would take the head out of its Service endpoints.
endpointPollInterval := 30 * time.Second
if intervalStr := os.Getenv("RAY_COLLECTOR_POLL_INTERVAL"); intervalStr != "" {
parsed, parseErr := time.ParseDuration(intervalStr)
if parseErr != nil {
logrus.Fatalf("Failed to parse RAY_COLLECTOR_POLL_INTERVAL: %v", parseErr)
if v := os.Getenv("RAY_COLLECTOR_POLL_INTERVAL"); v != "" {
parsed, err := time.ParseDuration(v)
if err == nil && parsed <= 0 {
err = errors.New("must be positive")
}
if parsed <= 0 {
logrus.Fatalf("RAY_COLLECTOR_POLL_INTERVAL must be positive, got: %s", intervalStr)
if err != nil {
logrus.Warnf("Invalid RAY_COLLECTOR_POLL_INTERVAL=%s (%v), using default %s", v, err, endpointPollInterval)
} else {
endpointPollInterval = parsed
}
endpointPollInterval = parsed
}
Comment thread
win5923 marked this conversation as resolved.

jsonData := make(map[string]interface{})
Expand Down Expand Up @@ -210,6 +217,14 @@ func main() {

sessionName := path.Base(activeSessionDir)

// The collector always runs as a sidecar in the Ray head pod, so the dashboard is
// reachable on localhost at Ray's default port. Only the head collector uses this.
// Override it when the dashboard listens on a non-default port.
dashboardAddress := "http://localhost:8265"
if v := os.Getenv("RAY_DASHBOARD_ADDRESS"); v != "" {
dashboardAddress = v
}

globalConfig := types.RayCollectorConfig{
RootDir: rayRootDir,
SessionDir: activeSessionDir,
Expand All @@ -219,7 +234,7 @@ func main() {
RayClusterNamespace: rayClusterNamespace,
PushInterval: pushInterval,
LogBatching: logBatching,
DashboardAddress: os.Getenv("RAY_DASHBOARD_ADDRESS"),
DashboardAddress: dashboardAddress,
OwnerKind: ownerKind,
OwnerName: ownerName,

Expand Down
126 changes: 126 additions & 0 deletions historyserver/config/ray-data.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,126 @@
apiVersion: ray.io/v1
kind: RayJob
metadata:
name: rayjob-ray-data
spec:
# Self-contained: this RayJob brings up its own cluster instead of attaching to an
# existing one via clusterSelector, so it also exercises the collector's shutdown path.
shutdownAfterJobFinishes: true
# Keeps the cluster alive long enough for at least one polling cycle after the job
# succeeds. Without it the cluster is deleted immediately and the datasets would only
# be captured by the collector's best-effort final poll during shutdown.
ttlSecondsAfterFinished: 30
entrypoint: |
python -c "
import ray
ray.init()

# materialize() is required: an unexecuted Dataset produces no stats, so
# /api/data/datasets/{job_id} would stay empty.
ds = ray.data.range(100).map_batches(lambda batch: batch).materialize()
print(f'Dataset rows: {ds.count()}')
"
rayClusterSpec:
# Head-only on purpose: worker collectors need the head Service FQDN in FQ_RAY_IP,
# which cannot be written here because KubeRay generates the cluster name.
headGroupSpec:
rayStartParams:
dashboard-host: 0.0.0.0
serviceType: ClusterIP
template:
spec:
containers:
- env:
- name: RAY_TMP_ROOT
value: &rayTmpRoot /tmp/ray
- name: RAY_enable_ray_event
value: "true"
- name: RAY_enable_core_worker_ray_event_to_aggregator
value: "true"
- name: RAY_DASHBOARD_AGGREGATOR_AGENT_EVENTS_EXPORT_ADDR
value: "http://localhost:8084/v1/events"
# in ray 2.52.0, we need to set RAY_DASHBOARD_AGGREGATOR_AGENT_EXPOSABLE_EVENT_TYPES
# in ray 2.53.0 (noy yet done). we need to set RAY_DASHBOARD_AGGREGATOR_AGENT_PUBLISHER_HTTP_ENDPOINT_EXPOSABLE_EVENT_TYPES
- name: RAY_DASHBOARD_AGGREGATOR_AGENT_EXPOSABLE_EVENT_TYPES
value: "TASK_DEFINITION_EVENT,TASK_LIFECYCLE_EVENT,ACTOR_TASK_DEFINITION_EVENT,
TASK_PROFILE_EVENT,DRIVER_JOB_DEFINITION_EVENT,DRIVER_JOB_LIFECYCLE_EVENT,
ACTOR_DEFINITION_EVENT,ACTOR_LIFECYCLE_EVENT,NODE_DEFINITION_EVENT,NODE_LIFECYCLE_EVENT"
Comment on lines +40 to +42

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's update the ray image version to 2.56.0 and use "ALL" here

image: rayproject/ray:2.52.0
imagePullPolicy: IfNotPresent
name: ray-head
securityContext:
allowPrivilegeEscalation: true
privileged: true
resources:
limits:
cpu: "5"
memory: 10G
requests:
cpu: "50m"
memory: 1G
volumeMounts:
- name: historyserver
mountPath: *rayTmpRoot
- name: collector
image: collector:v0.1.0
imagePullPolicy: IfNotPresent
env:
- name: POD_IP
valueFrom:
fieldRef:
fieldPath: status.podIP
# KubeRay generates the cluster name, so these are read back from the labels
# it stamps on the pod rather than hardcoded.
- name: RAY_CLUSTER_NAME
valueFrom:
fieldRef:
fieldPath: metadata.labels['ray.io/cluster']
- name: RAY_CLUSTER_NAMESPACE
valueFrom:
fieldRef:
fieldPath: metadata.namespace
# Hardcoded, unlike the cluster name: KubeRay puts ray.io/originated-from-*
# on the RayCluster but not on the pod, so the downward API cannot read them.
# Must match metadata.name above.
- name: OWNER_KIND
value: "RayJob"
- name: OWNER_NAME
value: "rayjob-ray-data"
# Only used to look up this pod's Ray NodeID, and the dashboard is in this
# same pod. A worker collector would need the head Service FQDN instead.
- name: FQ_RAY_IP
value: "localhost"
- name: RAY_TMP_ROOT
value: *rayTmpRoot
# Shorter than the 30s default so a full cycle fits inside
# ttlSecondsAfterFinished above.
- name: RAY_COLLECTOR_POLL_INTERVAL
value: "5s"
- name: S3DISABLE_SSL
value: "true"
- name: AWS_ACCESS_KEY_ID
value: minioadmin
- name: AWS_SECRET_ACCESS_KEY
value: minioadmin
- name: AWS_SESSION_TOKEN
value: ""
- name: S3_BUCKET
value: "ray-historyserver"
- name: S3_ENDPOINT
value: "minio-service.minio-dev:9000"
- name: S3_REGION
value: "test"
- name: S3FORCE_PATH_STYLE
value: "true"
command:
- collector
- --role=Head
- --runtime-class-name=s3
- --ray-root-dir=log
- --events-port=8084
volumeMounts:
- name: historyserver
mountPath: *rayTmpRoot
volumes:
- name: historyserver
emptyDir: {}
16 changes: 16 additions & 0 deletions historyserver/config/raycluster-azureblob.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,22 @@ spec:
value: raycluster-historyserver-head-svc.default.svc.cluster.local
- name: RAY_TMP_ROOT
value: *rayTmpRoot
# RAY_DASHBOARD_ADDRESS points the head collector at the Ray Dashboard in the same
# pod. Optional; defaults to http://localhost:8265. Uncomment only if the dashboard
# listens on a non-default port. Worker collectors do not use it.
# - name: RAY_DASHBOARD_ADDRESS
# value: "http://localhost:9265"
# RAY_COLLECTOR_POLL_INTERVAL sets how often the head collector polls the Ray
# Dashboard endpoints. Optional; defaults to 30s. Accepts Go duration format.
# - name: RAY_COLLECTOR_POLL_INTERVAL
# value: "1m"
# The head collector always polls its built-in endpoints (Serve applications,
# placement groups, and per-job Ray Data datasets). RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# is optional and adds more on top; uncomment to use it. Each comma-separated path
# must match what the dashboard frontend requests, query string included, because
# the storage key is derived from the request URI.
# - name: RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# value: "/nodes?view=summary"
Comment on lines +69 to +72

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm curious why the sample value is /nodes?view=summary. Only endpoints under /api/** that don't have dedicated handlers can be replayed.

ws.Path("/api").Consumes(restful.MIME_JSON).Produces(restful.MIME_JSON).Filter(RequestLogFilter) //.Filter(s.loginWrapper)

# reference: https://learn.microsoft.com/en-us/azure/storage/common/storage-use-azurite#connect-to-the-emulator-by-using-the-azure-storage-explorer
- name: AZURE_STORAGE_CONNECTION_STRING
value: "DefaultEndpointsProtocol=http;AccountName=devstoreaccount1;AccountKey=Eby8vdM02xNOcqFlqUwJPLlmEtlCDXJ1OUzFT50uSRZ6IFsuFq2UVErCz4I6tq/K1SZFPTOtr/KBHBeksoGMGw==;BlobEndpoint=http://azurite-service.azurite-dev.svc.cluster.local:10000/devstoreaccount1;"
Expand Down
37 changes: 16 additions & 21 deletions historyserver/config/raycluster-gcs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -63,27 +63,22 @@ spec:
value: *rayTmpRoot
- name: GCS_BUCKET
value: "${GCS_BUCKET}"
# RAY_DASHBOARD_ADDRESS is used by the head collector to fetch endpoints' results
# (e.g., /api/v0/cluster_metadata) from the Ray Dashboard running in the same pod.
# Only the head collector uses this; worker collectors do not need it.
# If your Ray Dashboard uses a non-default port (not 8265), update this value accordingly.
- name: RAY_DASHBOARD_ADDRESS
value: "http://localhost:8265"
# RAY_COLLECTOR_ADDITIONAL_ENDPOINTS is a comma-separated list of Ray Dashboard
# API endpoint paths that the head collector will periodically poll and store.
# Use this for endpoints whose data cannot be obtained via Ray events.
# You can add more endpoints, e.g., "/api/v0/placement_groups,/api/serve/applications/"
# Note: Only static endpoints (no dynamic path parameters like {job_id}) are supported.
- name: RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
value: "/api/v0/placement_groups?detail=1&limit=10000"
# Query params must match the Ray Dashboard frontend request exactly
# (see https://github.com/ray-project/ray/blob/cb9c80fee6a700efe61ea97987248ce82e3fa2e2/python/ray/dashboard/client/src/service/placementGroup.ts):
# detail=1 — include bundles and stats fields required by PlacementGroupTable
# limit=10000 — match the frontend's default limit to avoid truncation
# RAY_COLLECTOR_POLL_INTERVAL controls how often the collector polls the additional
# endpoints above. Accepts Go duration format (e.g., "30s", "1m", "5m").
- name: RAY_COLLECTOR_POLL_INTERVAL
value: "30s"
# RAY_DASHBOARD_ADDRESS points the head collector at the Ray Dashboard in the same
# pod. Optional; defaults to http://localhost:8265. Uncomment only if the dashboard
# listens on a non-default port. Worker collectors do not use it.
# - name: RAY_DASHBOARD_ADDRESS
# value: "http://localhost:9265"
# RAY_COLLECTOR_POLL_INTERVAL sets how often the head collector polls the Ray
# Dashboard endpoints. Optional; defaults to 30s. Accepts Go duration format.
# - name: RAY_COLLECTOR_POLL_INTERVAL
# value: "1m"
# The head collector always polls its built-in endpoints (Serve applications,
# placement groups, and per-job Ray Data datasets). RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# is optional and adds more on top; uncomment to use it. Each comma-separated path
# must match what the dashboard frontend requests, query string included, because
# the storage key is derived from the request URI.
# - name: RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# value: "/nodes?view=summary"
command:
- collector
- --role=Head
Expand Down
37 changes: 16 additions & 21 deletions historyserver/config/raycluster.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -66,27 +66,22 @@ spec:
value: raycluster-historyserver-head-svc.default.svc.cluster.local
- name: RAY_TMP_ROOT
value: *rayTmpRoot
# RAY_DASHBOARD_ADDRESS is used by the head collector to fetch endpoints' results
# (e.g., /api/v0/cluster_metadata) from the Ray Dashboard running in the same pod.
# Only the head collector uses this; worker collectors do not need it.
# If your Ray Dashboard uses a non-default port (not 8265), update this value accordingly.
- name: RAY_DASHBOARD_ADDRESS
value: "http://localhost:8265"
# RAY_COLLECTOR_ADDITIONAL_ENDPOINTS is a comma-separated list of Ray Dashboard
# API endpoint paths that the head collector will periodically poll and store.
# Use this for endpoints whose data cannot be obtained via Ray events.
# You can add more endpoints, e.g., "/api/v0/placement_groups,/api/serve/applications/"
# Note: Only static endpoints (no dynamic path parameters like {job_id}) are supported.
- name: RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# Query params must match the Ray Dashboard frontend request exactly
# (see https://github.com/ray-project/ray/blob/cb9c80fee6a700efe61ea97987248ce82e3fa2e2/python/ray/dashboard/client/src/service/placementGroup.ts):
# detail=1 — include bundles and stats fields required by PlacementGroupTable
# limit=10000 — match the frontend's default limit to avoid truncation
value: "/api/v0/placement_groups?detail=1&limit=10000"
# RAY_COLLECTOR_POLL_INTERVAL controls how often the collector polls the additional
# endpoints above. Accepts Go duration format (e.g., "30s", "1m", "5m").
- name: RAY_COLLECTOR_POLL_INTERVAL
value: "30s"
# RAY_DASHBOARD_ADDRESS points the head collector at the Ray Dashboard in the same
# pod. Optional; defaults to http://localhost:8265. Uncomment only if the dashboard
# listens on a non-default port. Worker collectors do not use it.
# - name: RAY_DASHBOARD_ADDRESS
# value: "http://localhost:9265"
# RAY_COLLECTOR_POLL_INTERVAL sets how often the head collector polls the Ray
# Dashboard endpoints. Optional; defaults to 30s. Accepts Go duration format.
# - name: RAY_COLLECTOR_POLL_INTERVAL
# value: "1m"
# The head collector always polls its built-in endpoints (Serve applications,
# placement groups, and per-job Ray Data datasets). RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# is optional and adds more on top; uncomment to use it. Each comma-separated path
# must match what the dashboard frontend requests, query string included, because
# the storage key is derived from the request URI.
# - name: RAY_COLLECTOR_ADDITIONAL_ENDPOINTS
# value: "/nodes?view=summary"
- name: S3DISABLE_SSL
value: "true"
- name: AWS_ACCESS_KEY_ID
Expand Down
Loading
Loading