Commit ef07ae87 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

chore: 清理遗留的Streamlit前端文件并添加SSH密钥配置脚本

parent 892e0326
Loading
Loading
Loading
Loading
+1 −1
Original line number Diff line number Diff line
@@ -9,7 +9,7 @@ docker-compose.override.yaml

# Conda
miniconda3/
*.sh
#*.sh

# Logs
logs/
+0 −335
Original line number Diff line number Diff line
# 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.
#

# Basic Airflow cluster configuration for CeleryExecutor with Redis and PostgreSQL.
#
# WARNING: This configuration is for local development. Do not use it in a production deployment.
#
# This configuration supports basic configuration using environment variables or an .env file
# The following variables are supported:
#
# AIRFLOW_IMAGE_NAME           - Docker image name used to run Airflow.
#                                Default: apache/airflow:3.0.5
# AIRFLOW_UID                  - User ID in Airflow containers
#                                Default: 50000
# AIRFLOW_PROJ_DIR             - Base path to which all the files will be volumed.
#                                Default: .
# Those configurations are useful mostly in case of standalone testing/running Airflow in test/try-out mode
#
# _AIRFLOW_WWW_USER_USERNAME   - Username for the administrator account (if requested).
#                                Default: airflow
# _AIRFLOW_WWW_USER_PASSWORD   - Password for the administrator account (if requested).
#                                Default: airflow
# _PIP_ADDITIONAL_REQUIREMENTS - Additional PIP requirements to add when starting all containers.
#                                Use this option ONLY for quick checks. Installing requirements at container
#                                startup is done EVERY TIME the service is started.
#                                A better way is to build a custom image or extend the official image
#                                as described in https://airflow.apache.org/docs/docker-stack/build.html.
#                                Default: ''
#
# Feel free to modify this file to suit your needs.
---
x-airflow-common:
  &airflow-common
  # In order to add custom dependencies or upgrade provider distributions you can use your extended image.
  # Comment the image line, place your Dockerfile in the directory where you placed the docker-compose.yaml
  # and uncomment the "build" line below, Then run `docker-compose build` to build the images.
  image: ${AIRFLOW_IMAGE_NAME:-apache/airflow:3.0.5}
  # build: .
  environment:
    &airflow-common-env
    AIRFLOW__CORE__EXECUTOR: CeleryExecutor
    AIRFLOW__CORE__AUTH_MANAGER: airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager
    AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow
    AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow
    AIRFLOW__CELERY__BROKER_URL: redis://:@redis:6379/0
    AIRFLOW__CORE__FERNET_KEY: ''
    AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'true'
    AIRFLOW__CORE__LOAD_EXAMPLES: 'true'
    AIRFLOW__CORE__EXECUTION_API_SERVER_URL: 'http://airflow-apiserver:8080/execution/'
    # yamllint disable rule:line-length
    # Use simple http server on scheduler for health checks
    # See https://airflow.apache.org/docs/apache-airflow/stable/administration-and-deployment/logging-monitoring/check-health.html#scheduler-health-check-server
    # yamllint enable rule:line-length
    AIRFLOW__SCHEDULER__ENABLE_HEALTH_CHECK: 'true'
    # WARNING: Use _PIP_ADDITIONAL_REQUIREMENTS option ONLY for a quick checks
    # for other purpose (development, test and especially production usage) build/extend Airflow image.
    _PIP_ADDITIONAL_REQUIREMENTS: ${_PIP_ADDITIONAL_REQUIREMENTS:-}
    # The following line can be used to set a custom config file, stored in the local config folder
    AIRFLOW_CONFIG: '/opt/airflow/config/airflow.cfg'
  volumes:
    - ${AIRFLOW_PROJ_DIR:-.}/dags:/opt/airflow/dags
    - ${AIRFLOW_PROJ_DIR:-.}/logs:/opt/airflow/logs
    - ${AIRFLOW_PROJ_DIR:-.}/config:/opt/airflow/config
    - ${AIRFLOW_PROJ_DIR:-.}/plugins:/opt/airflow/plugins
  user: "${AIRFLOW_UID:-50000}:0"
  depends_on:
    &airflow-common-depends-on
    redis:
      condition: service_healthy
    postgres:
      condition: service_healthy

services:
  postgres:
    image: postgres:13
    environment:
      POSTGRES_USER: airflow
      POSTGRES_PASSWORD: airflow
      POSTGRES_DB: airflow
    volumes:
      - postgres-db-volume:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD", "pg_isready", "-U", "airflow"]
      interval: 10s
      retries: 5
      start_period: 5s
    restart: always

  redis:
    # Redis is limited to 7.2-bookworm due to licencing change
    # https://redis.io/blog/redis-adopts-dual-source-available-licensing/
    image: redis:7.2-bookworm
    expose:
      - 6379
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 10s
      timeout: 30s
      retries: 50
      start_period: 30s
    restart: always

  airflow-apiserver:
    <<: *airflow-common
    command: api-server
    ports:
      - "8080:8080"
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:8080/api/v2/version"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-scheduler:
    <<: *airflow-common
    command: scheduler
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:8974/health"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-dag-processor:
    <<: *airflow-common
    command: dag-processor
    healthcheck:
      test: ["CMD-SHELL", 'airflow jobs check --job-type DagProcessorJob --hostname "$${HOSTNAME}"']
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-worker:
    <<: *airflow-common
    command: celery worker
    healthcheck:
      # yamllint disable rule:line-length
      test:
        - "CMD-SHELL"
        - 'celery --app airflow.providers.celery.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}" || celery --app airflow.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}"'
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    environment:
      <<: *airflow-common-env
      # Required to handle warm shutdown of the celery workers properly
      # See https://airflow.apache.org/docs/docker-stack/entrypoint.html#signal-propagation
      DUMB_INIT_SETSID: "0"
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-apiserver:
        condition: service_healthy
      airflow-init:
        condition: service_completed_successfully

  airflow-triggerer:
    <<: *airflow-common
    command: triggerer
    healthcheck:
      test: ["CMD-SHELL", 'airflow jobs check --job-type TriggererJob --hostname "$${HOSTNAME}"']
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

  airflow-init:
    <<: *airflow-common
    entrypoint: /bin/bash
    # yamllint disable rule:line-length
    command:
      - -c
      - |
        if [[ -z "${AIRFLOW_UID}" ]]; then
          echo
          echo -e "\033[1;33mWARNING!!!: AIRFLOW_UID not set!\e[0m"
          echo "If you are on Linux, you SHOULD follow the instructions below to set "
          echo "AIRFLOW_UID environment variable, otherwise files will be owned by root."
          echo "For other operating systems you can get rid of the warning with manually created .env file:"
          echo "    See: https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#setting-the-right-airflow-user"
          echo
          export AIRFLOW_UID=$$(id -u)
        fi
        one_meg=1048576
        mem_available=$$(($$(getconf _PHYS_PAGES) * $$(getconf PAGE_SIZE) / one_meg))
        cpus_available=$$(grep -cE 'cpu[0-9]+' /proc/stat)
        disk_available=$$(df / | tail -1 | awk '{print $$4}')
        warning_resources="false"
        if (( mem_available < 4000 )) ; then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough memory available for Docker.\e[0m"
          echo "At least 4GB of memory required. You have $$(numfmt --to iec $$((mem_available * one_meg)))"
          echo
          warning_resources="true"
        fi
        if (( cpus_available < 2 )); then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough CPUS available for Docker.\e[0m"
          echo "At least 2 CPUs recommended. You have $${cpus_available}"
          echo
          warning_resources="true"
        fi
        if (( disk_available < one_meg * 10 )); then
          echo
          echo -e "\033[1;33mWARNING!!!: Not enough Disk space available for Docker.\e[0m"
          echo "At least 10 GBs recommended. You have $$(numfmt --to iec $$((disk_available * 1024 )))"
          echo
          warning_resources="true"
        fi
        if [[ $${warning_resources} == "true" ]]; then
          echo
          echo -e "\033[1;33mWARNING!!!: You have not enough resources to run Airflow (see above)!\e[0m"
          echo "Please follow the instructions to increase amount of resources available:"
          echo "   https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#before-you-begin"
          echo
        fi
        echo
        echo "Creating missing opt dirs if missing:"
        echo
        mkdir -v -p /opt/airflow/{logs,dags,plugins,config}
        echo
        echo "Airflow version:"
        /entrypoint airflow version
        echo
        echo "Files in shared volumes:"
        echo
        ls -la /opt/airflow/{logs,dags,plugins,config}
        echo
        echo "Running airflow config list to create default config file if missing."
        echo
        /entrypoint airflow config list >/dev/null
        echo
        echo "Files in shared volumes:"
        echo
        ls -la /opt/airflow/{logs,dags,plugins,config}
        echo
        echo "Change ownership of files in /opt/airflow to ${AIRFLOW_UID}:0"
        echo
        chown -R "${AIRFLOW_UID}:0" /opt/airflow/
        echo
        echo "Change ownership of files in shared volumes to ${AIRFLOW_UID}:0"
        echo
        chown -v -R "${AIRFLOW_UID}:0" /opt/airflow/{logs,dags,plugins,config}
        echo
        echo "Files in shared volumes:"
        echo
        ls -la /opt/airflow/{logs,dags,plugins,config}

    # yamllint enable rule:line-length
    environment:
      <<: *airflow-common-env
      _AIRFLOW_DB_MIGRATE: 'true'
      _AIRFLOW_WWW_USER_CREATE: 'true'
      _AIRFLOW_WWW_USER_USERNAME: ${_AIRFLOW_WWW_USER_USERNAME:-airflow}
      _AIRFLOW_WWW_USER_PASSWORD: ${_AIRFLOW_WWW_USER_PASSWORD:-airflow}
      _PIP_ADDITIONAL_REQUIREMENTS: ''
    user: "0:0"

  airflow-cli:
    <<: *airflow-common
    profiles:
      - debug
    environment:
      <<: *airflow-common-env
      CONNECTION_CHECK_MAX_COUNT: "0"
    # Workaround for entrypoint issue. See: https://github.com/apache/airflow/issues/16252
    command:
      - bash
      - -c
      - airflow
    depends_on:
      <<: *airflow-common-depends-on

  # You can enable flower by adding "--profile flower" option e.g. docker-compose --profile flower up
  # or by explicitly targeted on the command line e.g. docker-compose up flower.
  # See: https://docs.docker.com/compose/profiles/
  flower:
    <<: *airflow-common
    command: celery flower
    profiles:
      - flower
    ports:
      - "5555:5555"
    healthcheck:
      test: ["CMD", "curl", "--fail", "http://localhost:5555/"]
      interval: 30s
      timeout: 10s
      retries: 5
      start_period: 30s
    restart: always
    depends_on:
      <<: *airflow-common-depends-on
      airflow-init:
        condition: service_completed_successfully

volumes:
  postgres-db-volume:
+0 −12
Original line number Diff line number Diff line
FROM python:3.9-slim

WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

EXPOSE 8501

CMD ["streamlit", "run", "app.py", "--server.port", "8501", "--server.address", "0.0.0.0", "--browser.gatherUsageStats", "false"]
+0 −184
Original line number Diff line number Diff line
import streamlit as st
import requests
import json
import pandas as pd
from datetime import datetime
from streamlit_autorefresh import st_autorefresh

import os

# 配置页面
st.set_page_config(
    page_title="CSST 任务统一控制台",
    page_icon="🚀",
    layout="wide",
    initial_sidebar_state="expanded"
)

# 后端 API 地址,支持通过环境变量传入
API_HOST = os.environ.get("API_GATEWAY_HOST", "api-gateway")
if not API_HOST:
    API_HOST = "api-gateway"
API_PORT = os.environ.get("API_GATEWAY_PORT", "38000")
API_BASE_URL = f"http://{API_HOST}:{API_PORT}/api/tasks"

GRAFANA_HOST = os.environ.get("MASTER_IP", "localhost")
if not GRAFANA_HOST:
    GRAFANA_HOST = "localhost"
GRAFANA_PORT = os.environ.get("GRAFANA_PORT", "33000")
GRAFANA_URL = f"http://{GRAFANA_HOST}:{GRAFANA_PORT}"

st.title("🚀 CSST 任务统一控制台 (Task Portal)")

# 侧边栏导航
page = st.sidebar.radio("导航", ["🔍 任务大厅", "📊 资源监控 (Grafana)"])

if page == "🔍 任务大厅":
    st.header("🔍 任务查询与管控")
    
    # 启用自动轮询刷新,每 5 秒刷新一次,最多刷新 10000 次
    st_autorefresh(interval=5000, limit=10000, key="task_hall_autorefresh")
    
    with st.expander("📝 智能条件筛选", expanded=True):
        st.markdown("通过动态提取业务字段进行筛选,或直接输入高级 JSON")
        
        # 尝试获取动态字段
        try:
            fields_res = requests.get(f"{API_BASE_URL}/metadata/fields", timeout=2)
            if fields_res.status_code == 200:
                fields_data = fields_res.json().get("fields", {})
            else:
                fields_data = {}
        except:
            fields_data = {}

        available_keys = list(fields_data.keys())

        col1, col2, col3 = st.columns([2, 2, 1])
        with col1:
            selected_field = st.selectbox("选择业务字段 (Key)", [""] + available_keys, help="这些字段是动态从最近的任务中提取的")
            if selected_field and fields_data.get(selected_field):
                available_values = fields_data[selected_field]
                exact_value = st.selectbox("匹配值 (Value)", [""] + available_values, help="该字段的常用取值")
            else:
                exact_value = st.text_input("匹配值 (Value)", help="精确匹配该字段的值")
        
        with col2:
            st.markdown("**(可选) 手动高级 JSON**")
            query_json_str = st.text_area(
                "JSON 格式", 
                value="{}",
                height=100
            )
            
        with col3:
            st.write("")
            st.write("")
            search_btn = st.button("开始检索 🔎", use_container_width=True)
            
    # 如果点击了搜索按钮,或者 session 里没有数据,或者这是通过自动刷新触发的重新渲染,都去后台拿最新数据
    if search_btn or 'search_results' not in st.session_state or st.session_state.get("task_hall_autorefresh", 0) > 0:
        try:
            # 组合条件
            query_json = json.loads(query_json_str)
            if selected_field and exact_value:
                query_json[selected_field] = exact_value
                
            response = requests.post(f"{API_BASE_URL}/search", json=query_json, params={"limit": 50})
            if response.status_code == 200:
                st.session_state.search_results = response.json()
            else:
                st.error(f"检索失败: {response.text}")
        except json.JSONDecodeError:
            st.error("JSON 格式错误,请检查!")
        except Exception as e:
            st.error(f"请求失败: {e}")

    if 'search_results' in st.session_state and st.session_state.search_results:
        results = st.session_state.search_results
        st.write(f"找到 **{len(results)}** 条匹配任务:")
        
        for task in results:
            task_id = task['task_id']
            status = task['status']
            dag_id = task['dag_id']
            
            # 状态颜色
            color_map = {
                "success": "green",
                "failed": "red",
                "cancelled": "orange",
                "running": "blue",
                "received": "gray",
                "retrying": "purple"
            }
            color = color_map.get(status, "black")
            
            with st.container():
                c1, c2, c3, c4 = st.columns([2, 1, 3, 2])
                with c1:
                    st.markdown(f"**Task ID**: `{task_id}`")
                    st.markdown(f"**DAG ID**: `{dag_id}`")
                with c2:
                    st.markdown(f"**状态**: :{color}[**{status.upper()}**]")
                with c3:
                    # 进度预估
                    if status not in ["success", "failed", "cancelled"]:
                        try:
                            eta_res = requests.get(f"{API_BASE_URL}/{task_id}/eta", timeout=2)
                            if eta_res.status_code == 200:
                                eta_data = eta_res.json()
                                progress = eta_data.get("progress_percent", 0)
                                eta_sec = eta_data.get("eta_seconds", 0)
                                st.progress(int(progress) / 100.0, text=f"预估进度: {progress}% (剩余约 {eta_sec} 秒)")
                            else:
                                st.write("进度加载中...")
                        except:
                            st.write("进度获取失败")
                    else:
                        st.progress(100, text="已结束")
                
                with c4:
                    bc1, bc2 = st.columns(2)
                    with bc1:
                        if st.button("❌ 取消", key=f"cancel_{task_id}"):
                            res = requests.post(f"{API_BASE_URL}/{task_id}/cancel")
                            if res.status_code == 200:
                                st.success("已发送取消请求!")
                            else:
                                st.error("取消失败")
                    with bc2:
                        if st.button("🔄 重试", key=f"retry_{task_id}"):
                            res = requests.post(f"{API_BASE_URL}/{task_id}/retry")
                            if res.status_code == 200:
                                st.success("已发送重试请求!")
                            else:
                                st.error("重试失败")
                st.divider()

elif page == "📊 资源监控 (Grafana)":
    st.header("📊 集群资源与队列监控")
    
    st.info(f"💡 **提示**: 如果下方的监控大盘由于浏览器跨域安全策略未能显示,请 [👉 点击这里直接在独立标签页打开 Grafana]({GRAFANA_URL})")
    
    st.markdown("""
        这里无缝嵌入了后端的 Grafana 看板。系统已为您预置了常用的监控大盘,请通过下方按钮切换视图:
    """)
    
    dashboard_choice = st.radio(
        "选择监控视图:",
        ["🖥️ 物理机资源 (Node Exporter)", "🐳 容器资源与运行状态 (cAdvisor)"],
        horizontal=True
    )
    
    if dashboard_choice == "🖥️ 物理机资源 (Node Exporter)":
        dash_url = f"{GRAFANA_URL}/d/node-exporter/node-exporter-full?orgId=1&kiosk=tv"
    else:
        dash_url = f"{GRAFANA_URL}/d/cadvisor-docker/cadvisor-docker?orgId=1&kiosk=tv"
    
    # 使用 iframe 嵌入 Grafana
    # 增加 allow="fullscreen" 属性,并通过 kiosk 模式隐藏 Grafana 侧边栏
    iframe_html = f"""
        <iframe src="{dash_url}" width="100%" height="800px" frameborder="0" allow="fullscreen"></iframe>
    """
    st.components.v1.html(iframe_html, height=800)
+0 −4
Original line number Diff line number Diff line
streamlit==1.33.0
requests==2.31.0
pandas==2.2.1
streamlit-autorefresh==1.0.1
 No newline at end of file
Loading