import { Button, Card, Space, Table, Tag, Toast, Typography } from '@douyinfe/semi-ui'; import { IconCopy } from '@douyinfe/semi-icons'; import { useEffect, useState } from 'react'; import { api } from '../api/client'; import type { CapacityMetrics, LinkHealth, OpsHealth } from '../api/types'; import { PageHeader } from '../components/PageHeader'; const statusColor: Record = { ok: 'green', warning: 'orange', error: 'red' }; function formatNumber(value?: number | null) { return value == null ? '-' : value.toLocaleString(); } function statusText(ok?: boolean) { return ok ? '正常' : '异常'; } function boolColor(ok?: boolean): 'green' | 'red' { return ok ? 'green' : 'red'; } function amapSecurityText(health: OpsHealth | null) { const runtime = health?.runtime; if (!runtime?.amapWebJsConfigured) return '地图未配置'; if (runtime.amapSecurityCodeExposed) return '安全码已暴露'; if (runtime.amapSecurityProxyEnabled) return '安全码未暴露'; return '安全代理未启用'; } function amapSecurityColor(health: OpsHealth | null): 'green' | 'orange' | 'red' { const runtime = health?.runtime; if (runtime?.amapSecurityCodeExposed) return 'red'; if (runtime?.amapWebJsConfigured && runtime?.amapSecurityProxyEnabled) return 'green'; return 'orange'; } type CapacityRisk = { key: string; title: string; value: string; category: string; severity: 'warning' | 'error'; finding: string; action: string; }; type CapacityMetricRow = { key: keyof CapacityMetrics; name: string; value: number; threshold: string; status: 'ok' | 'warning' | 'error'; detail: string; }; const emptyCapacityMetrics: CapacityMetrics = { activeConnections: 0, kafkaLag: 0, bridgeConsumerPending: 0, bridgeAckPending: 0, bridgeBatchPendingMessages: 0, fastWriterConsumerPending: 0, fastWriterAckPending: 0, fastWriterBatchPending: 0, historyBatchPending: 0, historyRowsPending: 0 }; function getCapacityMetrics(health: OpsHealth | null): CapacityMetrics { return { ...emptyCapacityMetrics, ...(health?.capacityMetrics ?? {}), activeConnections: health?.capacityMetrics?.activeConnections ?? health?.activeConnections ?? 0, kafkaLag: health?.capacityMetrics?.kafkaLag ?? health?.kafkaLag ?? 0 }; } function metricStatus(value: number, warningAt: number, errorAt?: number): 'ok' | 'warning' | 'error' { if (errorAt != null && value > errorAt) return 'error'; return value > warningAt ? 'warning' : 'ok'; } function capacityMetricRows(health: OpsHealth | null): CapacityMetricRow[] { const metrics = getCapacityMetrics(health); return [ { key: 'bridgeConsumerPending', name: 'NATS 桥接待消费', value: metrics.bridgeConsumerPending, threshold: '> 10,000 预警', status: metricStatus(metrics.bridgeConsumerPending, 10000), detail: 'RAW/FIELDS 从 NATS 进入 Kafka 前的待消费量。' }, { key: 'bridgeAckPending', name: 'NATS 桥接 ACK', value: metrics.bridgeAckPending, threshold: '> 100 严重', status: metricStatus(metrics.bridgeAckPending, 0, 100), detail: '桥接消费者已拉取但未确认的消息,持续不为 0 需要优先排查。' }, { key: 'fastWriterConsumerPending', name: '快速写入待消费', value: metrics.fastWriterConsumerPending, threshold: '> 10,000 预警', status: metricStatus(metrics.fastWriterConsumerPending, 10000), detail: '实时数据写 TDengine、Redis 前的 NATS 待消费量。' }, { key: 'fastWriterAckPending', name: '快速写入 ACK', value: metrics.fastWriterAckPending, threshold: '> 10 严重', status: metricStatus(metrics.fastWriterAckPending, 0, 10), detail: '快速写入消费者已拉取但未确认的消息。' }, { key: 'historyBatchPending', name: '历史批次待写', value: metrics.historyBatchPending, threshold: '> 0 预警', status: metricStatus(metrics.historyBatchPending, 0), detail: '历史落库批处理内存队列,正常应快速归零。' }, { key: 'historyRowsPending', name: '历史行待写', value: metrics.historyRowsPending, threshold: '> 0 预警', status: metricStatus(metrics.historyRowsPending, 0), detail: '历史落库批处理待写行数,用于判断 TDengine/MySQL 写入是否滞后。' }, { key: 'kafkaLag', name: 'Kafka 消费 Lag', value: metrics.kafkaLag, threshold: '> 0 预警', status: metricStatus(metrics.kafkaLag, 0), detail: '后续统计和历史消费链路的 Kafka 积压。' }, { key: 'activeConnections', name: '网关活跃连接', value: metrics.activeConnections, threshold: '> 100,000 预警', status: metricStatus(metrics.activeConnections, 100000), detail: '32960、808、MQTT 当前活跃连接总量。' } ]; } function firstNumber(value: string) { const match = value.match(/[\d,.]+/); return match ? match[0] : '-'; } function capacityRisk(finding: string): CapacityRisk { const lower = finding.toLowerCase(); if (lower.includes('bridge') && lower.includes('ack pending')) { return { key: finding, title: 'NATS 桥接 ACK 未确认', value: firstNumber(finding), category: 'NATS -> Kafka', severity: 'error', finding, action: '优先检查 Kafka 是否监听 9092、磁盘是否打满,以及桥接服务是否持续写入成功。' }; } if (lower.includes('bridge') && lower.includes('consumer pending')) { return { key: finding, title: 'NATS 桥接待消费', value: firstNumber(finding), category: 'NATS -> Kafka', severity: 'warning', finding, action: '确认 Kafka 已恢复后观察 pending 是否持续下降;若不下降,扩容桥接或降低 Kafka topic 保留压力。' }; } if (lower.includes('kafka') || lower.includes('lag')) { return { key: finding, title: 'Kafka 消费积压', value: firstNumber(finding), category: 'Kafka', severity: 'warning', finding, action: '检查消费者组 lag、topic retention、broker 磁盘和生产写入错误。' }; } if (lower.includes('redis')) { return { key: finding, title: 'Redis 在线投影异常', value: firstNumber(finding), category: 'Redis', severity: 'warning', finding, action: '核对 realtime KV、online TTL 和协议实时帧过滤,确认车辆在线统计来源。' }; } if (lower.includes('connection')) { return { key: finding, title: '网关连接压力', value: firstNumber(finding), category: 'Gateway', severity: 'warning', finding, action: '查看 32960、808、MQTT 连接数、系统 FD、CPU 和接入端口队列。' }; } return { key: finding, title: '容量风险', value: firstNumber(finding), category: 'General', severity: 'warning', finding, action: '进入对应链路日志和指标继续定位,确认是否影响车辆实时服务。' }; } function opsCapacityHandoffText(health: OpsHealth | null) { const metrics = getCapacityMetrics(health); const risks = (health?.capacityFindings ?? []).map(capacityRisk); const unhealthyLinks = (health?.linkHealth ?? []).filter((item) => item.status !== 'ok'); const runtime = health?.runtime; return [ '【车辆数据中台运维容量交接】', `运行版本:${runtime?.platformRelease || '-'}`, `Kafka Lag:${formatNumber(metrics.kafkaLag)}`, `活跃连接:${formatNumber(metrics.activeConnections)}`, `Redis 在线 Key:${formatNumber(health?.redisOnlineKeys)}`, `NATS桥接:consumer pending ${formatNumber(metrics.bridgeConsumerPending)} / ack pending ${formatNumber(metrics.bridgeAckPending)} / batch pending ${formatNumber(metrics.bridgeBatchPendingMessages)}`, `快速写入:consumer pending ${formatNumber(metrics.fastWriterConsumerPending)} / ack pending ${formatNumber(metrics.fastWriterAckPending)} / batch pending ${formatNumber(metrics.fastWriterBatchPending)}`, `历史写入:batch pending ${formatNumber(metrics.historyBatchPending)} / rows pending ${formatNumber(metrics.historyRowsPending)}`, `存储状态:TDengine ${statusText(health?.tdengineWritable)} / MySQL ${statusText(health?.mysqlWritable)}`, `地图配置:Web JS ${runtime?.amapWebJsConfigured ? '已配置' : '未配置'} / API ${runtime?.amapApiConfigured ? '已配置' : '未配置'} / ${amapSecurityText(health)}`, `异常链路:${unhealthyLinks.length > 0 ? unhealthyLinks.map((item) => `${item.name}=${item.status}(${item.detail || '-'})`).join(';') : '正常'}`, '', '容量发现:', ...(risks.length > 0 ? risks.map((risk, index) => `${index + 1}. [${risk.category}] ${risk.finding} -> ${risk.action}`) : ['暂无容量风险']), '', '处置顺序:', '1. 先看 Kafka Lag 与 NATS ACK pending,确认消息是否卡在桥接或下游消费', '2. 再看快速写入 pending,确认 Redis/TDengine 实时链路是否影响在线和实时数据', '3. 再看历史写入 pending,确认轨迹、RAW、统计是否滞后', '4. 最后核对 MySQL、TDengine、Redis 和地图配置,确认中台查询体验是否受影响' ].join('\n'); } async function copyText(value: string, label: string) { try { await navigator.clipboard.writeText(value); Toast.success(`已复制${label}`); } catch { Toast.error(`复制${label}失败`); } } export function OpsQuality() { const [health, setHealth] = useState(null); const [loading, setLoading] = useState(true); const load = () => { setLoading(true); api.opsHealth() .then(setHealth) .catch((error: Error) => Toast.error(error.message)) .finally(() => setLoading(false)); }; useEffect(() => { load(); }, []); const runtime = health?.runtime; const capacityRisks = (health?.capacityFindings ?? []).map(capacityRisk); const capacityRows = capacityMetricRows(health); const bridgeRisks = capacityRisks.filter((item) => item.category === 'NATS -> Kafka'); return (
)} />
{formatNumber(health?.kafkaLag)}
Kafka Lag
{formatNumber(health?.activeConnections)}
活跃连接
{formatNumber(health?.redisOnlineKeys)}
Redis 在线 Key
{runtime?.platformRelease || '-'}
运行版本
TDengine 写入
{statusText(health?.tdengineWritable)}
时序历史、RAW 和统计依赖 TDengine 可写状态。
MySQL 写入
{statusText(health?.mysqlWritable)}
车辆身份、绑定、服务配置和中台查询依赖 MySQL。
请求超时
{runtime?.requestTimeoutMs ? `${runtime.requestTimeoutMs.toLocaleString()} ms` : '-'}
接口层读取 Redis、TDengine、MySQL 的统一超时保护。
{runtime?.amapWebJsConfigured ? 'Web JS Key 已配置' : 'Web JS Key 未配置'} {runtime?.amapApiConfigured ? '服务端 API Key 已配置' : '服务端 API Key 未配置'} {amapSecurityText(health)} {runtime?.amapSecurityServiceHost || '代理未启用'} pagination={false} dataSource={capacityRows} rowKey={(row?: CapacityMetricRow) => row?.key ?? ''} columns={[ { title: '指标', dataIndex: 'name' }, { title: '当前值', width: 140, render: (_: unknown, row: CapacityMetricRow) => ( {formatNumber(row.value)} ) }, { title: '状态', width: 120, render: (_: unknown, row: CapacityMetricRow) => {row.status} }, { title: '阈值', width: 150, dataIndex: 'threshold' }, { title: '说明', dataIndex: 'detail' } ]} /> {capacityRisks.length ? (
{capacityRisks.map((item) => (
{item.category} {item.finding}
{item.value}
{item.title} {item.action}
))}
{bridgeRisks.length ? (
桥接处置顺序 先确认 Kafka broker 监听和磁盘空间,再观察 ACK pending 是否回到 0,最后确认 consumer pending 持续下降。
) : null}
) : null} pagination={false} dataSource={health?.linkHealth ?? []} rowKey={(row?: LinkHealth) => `${row?.name ?? ''}-${row?.status ?? ''}`} columns={[ { title: '链路', dataIndex: 'name' }, { title: '状态', width: 120, render: (_: unknown, row: LinkHealth) => {row.status} }, { title: '说明', dataIndex: 'detail' } ]} />
); }