From 7b651c95a9b3e9660c8c568d9e3926f62f4a5d39 Mon Sep 17 00:00:00 2001 From: Cristhian Garcia Date: Wed, 25 Mar 2026 08:00:10 -0500 Subject: [PATCH] feat: add vector aggregator --- tutoraspects/patches/k8s-deployments | 100 ++++++++++++++++-- tutoraspects/patches/k8s-jobs | 2 +- .../patches/kustomization-configmapgenerator | 13 ++- tutoraspects/plugin.py | 6 +- .../templates/aspects/apps/vector/file.toml | 15 --- .../aspects/apps/vector/k8s-agent.toml | 19 ++++ .../aspects/apps/vector/k8s-aggregator.toml | 8 ++ .../templates/aspects/apps/vector/k8s.toml | 13 --- .../apps/vector/partials/common-post.toml | 37 ++----- 9 files changed, 141 insertions(+), 72 deletions(-) delete mode 100644 tutoraspects/templates/aspects/apps/vector/file.toml create mode 100644 tutoraspects/templates/aspects/apps/vector/k8s-agent.toml create mode 100644 tutoraspects/templates/aspects/apps/vector/k8s-aggregator.toml delete mode 100644 tutoraspects/templates/aspects/apps/vector/k8s.toml diff --git a/tutoraspects/patches/k8s-deployments b/tutoraspects/patches/k8s-deployments index 8d7a62047..838d1c5b3 100644 --- a/tutoraspects/patches/k8s-deployments +++ b/tutoraspects/patches/k8s-deployments @@ -446,7 +446,7 @@ kind: ServiceAccount metadata: name: vector labels: - app.kubernetes.io/name: vector + app.kubernetes.io/name: vector-agent automountServiceAccountToken: true --- apiVersion: rbac.authorization.k8s.io/v1 @@ -469,7 +469,7 @@ kind: ClusterRoleBinding metadata: name: vector labels: - app.kubernetes.io/name: vector + app.kubernetes.io/name: vector-agent roleRef: apiGroup: rbac.authorization.k8s.io kind: ClusterRole @@ -489,18 +489,18 @@ description: "This priority class should be used for Vector service pods only." apiVersion: apps/v1 kind: DaemonSet metadata: - name: vector + name: vector-agent labels: - app.kubernetes.io/name: vector + app.kubernetes.io/name: vector-agent spec: selector: matchLabels: - app.kubernetes.io/name: vector + app.kubernetes.io/name: vector-agent minReadySeconds: 0 template: metadata: labels: - app.kubernetes.io/name: vector + app.kubernetes.io/name: vector-agent vector.dev/exclude: "true" spec: serviceAccountName: vector @@ -508,6 +508,8 @@ spec: containers: - name: vector image: {{ DOCKER_IMAGE_VECTOR }} + command: ["vector"] + args: ["--config", "/etc/vector/vector.toml"] env: - name: VECTOR_SELF_NODE_NAME valueFrom: @@ -534,7 +536,10 @@ spec: mountPath: /var/log/ - mountPath: /etc/vector/vector.toml name: config - subPath: k8s.toml + subPath: k8s-agent.toml + readOnly: true + - name: varlibdockercontainers + mountPath: /var/lib/docker/containers readOnly: true securityContext: allowPrivilegeEscalation: false @@ -543,13 +548,92 @@ spec: volumes: - name: config configMap: - name: vector-config + name: vector-agent-config - name: data hostPath: path: /var/lib/vector - name: var-log hostPath: path: /var/log/ + - name: varlibdockercontainers + hostPath: + path: /var/lib/docker/containers +--- +# Vector Aggregator +# https://vector.dev/docs/setup/going-to-prod/arch/aggregator/ +--- +apiVersion: v1 +kind: Service +metadata: + name: vector-aggregator + labels: + app.kubernetes.io/name: vector-aggregator +spec: + clusterIP: None + selector: + app.kubernetes.io/name: vector-aggregator + ports: + - name: vector + port: {{ ASPECTS_VECTOR_AGGREGATOR_PORT }} + targetPort: vector +--- +apiVersion: apps/v1 +kind: StatefulSet +metadata: + name: vector-aggregator + labels: + app.kubernetes.io/name: vector-aggregator +spec: + serviceName: vector-aggregator + replicas: {{ ASPECTS_VECTOR_AGGREGATOR_REPLICAS }} + selector: + matchLabels: + app.kubernetes.io/name: vector-aggregator + template: + metadata: + labels: + app.kubernetes.io/name: vector-aggregator + spec: + containers: + - name: vector + image: {{ DOCKER_IMAGE_VECTOR }} + command: ["vector"] + args: ["--config", "/etc/vector/vector.toml"] + env: + - name: VECTOR_SELF_POD_NAME + valueFrom: + fieldRef: + fieldPath: metadata.name + - name: VECTOR_SELF_POD_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace + - name: VECTOR_LOG + value: info + ports: + - name: vector + containerPort: {{ ASPECTS_VECTOR_AGGREGATOR_PORT }} + volumeMounts: + - name: vector-data + mountPath: /vector-data-dir + - mountPath: /etc/vector/vector.toml + name: config + subPath: k8s-aggregator.toml + readOnly: true + securityContext: + allowPrivilegeEscalation: false + volumes: + - name: config + configMap: + name: vector-aggregator-config + volumeClaimTemplates: + - metadata: + name: vector-data + spec: + accessModes: ["ReadWriteOnce"] + resources: + requests: + storage: "{{ ASPECTS_VECTOR_AGGREGATOR_STORAGE_SIZE }}" {% endif %} {% if ASPECTS_ENABLE_EVENT_BUS_CONSUMER %} diff --git a/tutoraspects/patches/k8s-jobs b/tutoraspects/patches/k8s-jobs index 5467a8a27..cbc697149 100644 --- a/tutoraspects/patches/k8s-jobs +++ b/tutoraspects/patches/k8s-jobs @@ -10,7 +10,7 @@ spec: spec: restartPolicy: Never containers: - - name: aspects + - name: aspects-job env: - name: VENV_DIR value: /opt/venv diff --git a/tutoraspects/patches/kustomization-configmapgenerator b/tutoraspects/patches/kustomization-configmapgenerator index a70f8abd5..dde00475d 100644 --- a/tutoraspects/patches/kustomization-configmapgenerator +++ b/tutoraspects/patches/kustomization-configmapgenerator @@ -103,7 +103,16 @@ {% endif %} {% if RUN_VECTOR %} -- name: vector-config +- name: vector-agent-config files: - - plugins/aspects/apps/vector/k8s.toml + - plugins/aspects/apps/vector/k8s-agent.toml + options: + labels: + app.kubernetes.io/name: vector-agent +- name: vector-aggregator-config + files: + - plugins/aspects/apps/vector/k8s-aggregator.toml + options: + labels: + app.kubernetes.io/name: vector-aggregator {% endif %} diff --git a/tutoraspects/plugin.py b/tutoraspects/plugin.py index ff6a5437a..ae0414281 100644 --- a/tutoraspects/plugin.py +++ b/tutoraspects/plugin.py @@ -31,8 +31,6 @@ # Each new setting is a pair: (setting_name, default_value). # Prefix your setting names with 'ASPECTS_'. ("ASPECTS_VERSION", __version__), - # For our default deployment we currently use Celery -> Ralph for transport, - # so Vector is off by default. ("RUN_VECTOR", True), ("RUN_CLICKHOUSE", True), ("RUN_RALPH", False), @@ -186,6 +184,10 @@ ("ASPECTS_XAPI_S3_SINK_TIMEOUT_SECS", "600"), ("ASPECTS_VECTOR_DATABASE", "openedx"), ("ASPECTS_VECTOR_RAW_TRACKING_LOGS_TABLE", "_tracking"), + ("ASPECTS_VECTOR_AGGREGATOR_PORT", "6000"), + ("ASPECTS_VECTOR_AGGREGATOR_REPLICAS", 1), + ("ASPECTS_VECTOR_AGGREGATOR_BUFFER_MAX_SIZE", "1073741824"), + ("ASPECTS_VECTOR_AGGREGATOR_STORAGE_SIZE", "2Gi"), ("ASPECTS_DATA_TTL_EXPRESSION", "toDateTime(emission_time) + INTERVAL 1 YEAR"), ("ASPECTS_ALEMBIC_MIGRATIONS_DATABASE", "{{RALPH_DATABASE}}"), # Make sure LMS / CMS have event-routing-backends installed diff --git a/tutoraspects/templates/aspects/apps/vector/file.toml b/tutoraspects/templates/aspects/apps/vector/file.toml deleted file mode 100644 index 87aa5ede4..000000000 --- a/tutoraspects/templates/aspects/apps/vector/file.toml +++ /dev/null @@ -1,15 +0,0 @@ -{% include "aspects/apps/vector/partials/common-pre.toml" %} - -### Sources -# Capture logs from tracking.log -[sources.tracking_log_file] -type = "file" -include = ["/var/log/openedx/tracking.log"] - -[transforms.openedx_containers] -type = "filter" -# no-op filter: created for future-proof compatibility -condition = "true" -inputs = ["tracking_log_file"] - -{% include "aspects/apps/vector/partials/common-post.toml" %} diff --git a/tutoraspects/templates/aspects/apps/vector/k8s-agent.toml b/tutoraspects/templates/aspects/apps/vector/k8s-agent.toml new file mode 100644 index 000000000..049c7982f --- /dev/null +++ b/tutoraspects/templates/aspects/apps/vector/k8s-agent.toml @@ -0,0 +1,19 @@ +{% include "aspects/apps/vector/partials/common-pre.toml" %} + +### Sources +# Capture logs from kubernetes +[sources.kubernetes_logs] +type = "kubernetes_logs" +extra_namespace_label_selector = "kubernetes.io/metadata.name={{ K8S_NAMESPACE }}" +glob_minimum_cooldown_ms = 2000 + +[transforms.openedx_containers] +type = "filter" +inputs = ["kubernetes_logs"] +condition = 'includes(["lms", "cms", "lms-worker", "cms-worker", "lms-job", "cms-job", "aspects-job", "aspects-consumer"], .kubernetes.container_name)' + +[sinks.to_aggregator] +type = "vector" +inputs = ["openedx_containers"] +address = "vector-aggregator.{{ K8S_NAMESPACE }}:{{ ASPECTS_VECTOR_AGGREGATOR_PORT }}" +compression = true diff --git a/tutoraspects/templates/aspects/apps/vector/k8s-aggregator.toml b/tutoraspects/templates/aspects/apps/vector/k8s-aggregator.toml new file mode 100644 index 000000000..d318169ff --- /dev/null +++ b/tutoraspects/templates/aspects/apps/vector/k8s-aggregator.toml @@ -0,0 +1,8 @@ +{% include "aspects/apps/vector/partials/common-pre.toml" %} + +### Sources +[sources.openedx_containers] +type = "vector" +address = "0.0.0.0:{{ ASPECTS_VECTOR_AGGREGATOR_PORT }}" + +{% include "aspects/apps/vector/partials/common-post.toml" %} diff --git a/tutoraspects/templates/aspects/apps/vector/k8s.toml b/tutoraspects/templates/aspects/apps/vector/k8s.toml deleted file mode 100644 index 8af523232..000000000 --- a/tutoraspects/templates/aspects/apps/vector/k8s.toml +++ /dev/null @@ -1,13 +0,0 @@ -{% include "aspects/apps/vector/partials/common-pre.toml" %} - -### Sources -# Capture logs from kubernetes -[sources.kubernetes_logs] -type = "kubernetes_logs" -extra_namespace_label_selector = "kubernetes.io/metadata.name={{ K8S_NAMESPACE }}" -[transforms.openedx_containers] -type = "filter" -inputs = ["kubernetes_logs"] -condition = '.kubernetes.pod_namespace == "{{ K8S_NAMESPACE }}" && includes(["lms", "cms", "lms-worker", "cms-worker", "lms-job", "cms-job", "aspects-job", "aspects-consumer"], .kubernetes.container_name)' - -{% include "aspects/apps/vector/partials/common-post.toml" %} diff --git a/tutoraspects/templates/aspects/apps/vector/partials/common-post.toml b/tutoraspects/templates/aspects/apps/vector/partials/common-post.toml index 35366ed48..8fa6a7d2f 100644 --- a/tutoraspects/templates/aspects/apps/vector/partials/common-post.toml +++ b/tutoraspects/templates/aspects/apps/vector/partials/common-post.toml @@ -31,22 +31,6 @@ if err_timestamp != null { drop_on_error = true drop_on_abort = true - -[transforms.tracking_debug] -type = "remap" -inputs = ["tracking"] -# Time formats: https://docs.rs/chrono/0.4.19/chrono/format/strftime/index.html#specifiers -source = ''' -.message = parse_json!(.message) -''' - -# Log all events to stdout, for debugging -[sinks.out] -type = "console" -inputs = ["tracking_debug"] -encoding.codec = "json" -encoding.only_fields = ["time", "message.context.course_id", "message.context.user_id", "message.name"] - # # Send logs to clickhouse [sinks.clickhouse] type = "clickhouse" @@ -70,7 +54,6 @@ healthcheck = true type = "remap" inputs = ["openedx_containers"] # Time formats: https://docs.rs/chrono/0.4.19/chrono/format/strftime/index.html#specifiers - source = ''' parsed, err_regex = parse_regex(.message, r'^.* \[xapi_tracking\] [^{}]* (?P\{.*\})$') if err_regex != null { @@ -99,20 +82,6 @@ event_id = parsed_json.id drop_on_error = true drop_on_abort = true -[transforms.xapi_debug] -type = "remap" -inputs = ["xapi"] -# Time formats: https://docs.rs/chrono/0.4.19/chrono/format/strftime/index.html#specifiers -source = ''' -.message = parse_json!(.event) -''' - -[sinks.out_xapi] -type = "console" -inputs = ["xapi_debug"] -encoding.codec = "json" -encoding.only_fields = ["event_id", "emission_time", "event"] - [sinks.clickhouse_xapi] type = "clickhouse" auth.strategy = "basic" @@ -127,6 +96,9 @@ endpoint = "{% if CLICKHOUSE_SECURE_CONNECTION %}https{% else %}http{% endif %}: database = "{{ ASPECTS_VECTOR_DATABASE }}" table = "{{ ASPECTS_RAW_XAPI_TABLE }}" healthcheck = false +buffer.type = "disk" +buffer.max_size = {{ ASPECTS_VECTOR_AGGREGATOR_BUFFER_MAX_SIZE }} +buffer.when_full = "block" {% if ASPECTS_XAPI_S3_BUCKET %} [sinks.s3_xapi] @@ -147,6 +119,9 @@ compression = "zstd" batch.max_events = {{ ASPECTS_XAPI_S3_SINK_MAX_EVENTS }} batch.timeout_secs = {{ ASPECTS_XAPI_S3_SINK_TIMEOUT_SECS }} framing.method = "newline_delimited" +buffer.type = "disk" +buffer.max_size = {{ ASPECTS_VECTOR_AGGREGATOR_BUFFER_MAX_SIZE }} +buffer.when_full = "block" {% endif %}