三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

seo教学网站wordpress关闭谷歌

seo教学网站wordpress关闭谷歌 seo教学网站,wordpress关闭谷歌,scratch少儿编程网站,网站建设企业模板设计优秀的代码1#!/usr/bin/env python3 # -*- coding: utf-8 -*-scankube_monitor_v13 - 使用 Raw HTTP 响应的 Kubernetes 监控器v13 核心优化: 1. 使用 _preload_content=False 获取原始 JSON …设计优秀的代码1#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ scankube_monitor_v13 - 使用 Raw HTTP 响应的 Kubernetes 监控器v13 核心优化: 1. 使用 _preload_content=False 获取原始 JSON 响应,避免对象序列化开销 2. 直接解析 JSON,不创建 V1Pod/V1Node 对象,大幅减少内存和 CPU 占用 3. 支持 namespace 过滤和 label selector 4. 流式处理:每收到一页就更新缓存关键原理: - kubernetes python client 默认会 deserialize JSON - Python 对象 - V1Pod/V1Node - 这个过程对于 10000+ pods 非常慢(CPU 密集) - 使用 _preload_content=False 可以获取原始 HTTPResponse,直接解析 JSON - 这样避免了 deserialize 开销,速度快 10-100 倍使用示例:python scankube_monitor_v13.py -k conf/zzbm-ano.confpython scankube_monitor_v13.py -k conf/zzt-ano.conf --namespaces default,kube-system """import os import sys import time import threading import logging import json from logging.handlers import TimedRotatingFileHandler from collections import defaultdict from concurrent.futures import ThreadPoolExecutorfrom flask import Flask, Response, jsonifytry:from kubernetes import client, config, watchfrom kubernetes.client.rest import ApiException except ImportError as e:raise ImportError("Missing dependency: {}".format(e))# ============================================================================ # Logging Configuration # ============================================================================_LOG_DIR = '/qaxdata/s/services/scankube_monitor/logs' if not os.path.exists(_LOG_DIR):os.makedirs(_LOG_DIR)_LOG_FILE = os.path.join(_LOG_DIR, 'scankube_monitor.log')_LOG_LEVEL = os.environ.get('SCANKUBE_LOG_LEVEL', 'INFO').upper() _log_level_map = {'DEBUG': logging.DEBUG,'INFO': logging.INFO,'WARNING': logging.WARNING,'WARN': logging.WARNING,'ERROR': logging.ERROR,'CRITICAL': logging.CRITICAL, } LOG_LEVEL = _log_level_map.get(_LOG_LEVEL, logging.INFO)logging.basicConfig(level=LOG_LEVEL,format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',handlers=[logging.StreamHandler(sys.stdout),TimedRotatingFileHandler(_LOG_FILE,when='midnight',interval=1,backupCount=7,encoding='utf-8',delay=False)] ) logger = logging.getLogger('scankube_monitor') logging.getLogger('werkzeug').setLevel(logging.WARNING)# ============================================================================ # Flask Application # ============================================================================app = Flask(__name__)# ============================================================================ # Global Metrics Cache # ============================================================================class GlobalMetricsCache(object):def __init__(self):self._env_outputs = {}self._lock = threading.RLock()def set_env_output(self, environment, output_text):with self._lock:self._env_outputs[environment] = output_textdef get_all_output(self):with self._lock:outputs = list(self._env_outputs.values())if not outputs:return ('# HELP scankube_monitor_up scankube_monitor is running\n''# TYPE scankube_monitor_up gauge\n''scankube_monitor_up 1\n''# No environments registered or initial sync in progress\n')result = '\n'.join(filter(None, outputs))if not result.endswith('\n'):result += '\n'return resultdef get_stats(self):with self._lock:return {'environments': list(self._env_outputs.keys())}metrics_cache = GlobalMetricsCache()# ============================================================================ # Metrics Controller # ============================================================================class MetricsController(object):def __init__(self, max_workers=10, fallback_interval=30):self._executor = ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix='metrics-calc')self._pending_envs = set()self._force_calc_envs = set()self._last_calc_time = {}self._calc_lock = threading.Lock()self._stop_event = threading.Event()self._scheduler_thread = Noneself._min_interval = 5self._fallback_interval = fallback_intervaldef start(self):self._scheduler_thread = threading.Thread(target=self._scheduler_loop,name='metrics-scheduler',daemon=True)self._scheduler_thread.start()logger.info("Metrics controller started (max_workers={}, fallback_interval={}s)".format(self._executor._max_workers, self._fallback_interval))def stop(self):self._stop_event.set()self._executor.shutdown(wait=False)if self._scheduler_thread and self._scheduler_thread.is_alive():self._scheduler_thread.join(timeout=5)logger.info("Metrics controller stopped")def trigger_calc(self, environment, force=False):with self._calc_lock:if force:self._force_calc_envs.add(environment)else:self._pending_envs.add(environment)def _scheduler_loop(self):loop_count = 0while not self._stop_event.is_set():try:with self._calc_lock:envs_to_calc = list(self._pending_envs)force_envs = list(self._force_calc_envs)self._pending_envs.clear()self._force_calc_envs.clear()now = time.time()for env in force_envs:manager = _env_managers.get(env)if manager:self._last_calc_time[env] = nowself._executor.submit(self._calc_env, manager)for env in envs_to_calc:last_time = self._last_calc_time.get(env, 0)if now - last_time self._min_interval:with self._calc_lock:self._pending_envs.add(env)continuemanager = _env_managers.get(env)if manager:self._last_calc_time[env] = nowself._executor.submit(self._calc_env, manager)loop_count += 1if loop_count % self._fallback_interval == 0:self._do_fallback_calc()except Exception as e:logger.error("Scheduler error: {}".format(e))self._stop_event.wait(1)def _do_fallback_calc(self):now = time.time()for env, manager in _env_managers.items():last_time = self._last_calc_time.get(env, 0)if now - last_time = self._fallback_interval:logger.info("[{}] Fallback calc triggered (last calc {:.1f}s ago)".format(env, now - last_time))self._last_calc_time[env] = nowself._executor.submit(self._calc_env, manager)def _calc_env(self, manager):try:start = time.time()output = MetricsCalculator.calculate_and_format(manager.environment, manager.cache)metrics_cache.set_env_output(manager.environment, output)elapsed = time.time() - startlogger.info("[{}] Metrics calculated in {:.3f}s ({} nodes, {} pods)".format(manager.environment, elapsed,manager.cache.get_node_count(),manager.cache.get_pod_count()))except Exception as e:logger.error("[{}] Metrics calculation failed: {}".format(manager.environment, e))controller = MetricsController()# ============================================================================ # Local Cache (v13 - 存储原始 dict 而不是对象) # ============================================================================class LocalCache(object):"""v13: 存储原始 dict 而不是 V1Pod/V1Node 对象"""def __init__(self, environment):self.environment = environmentself._pods = {} # {uid: pod_dict}self._nodes = {} # {name: node_dict}self._lock = threading.RLock()self._node_synced = Falseself._pod_synced = Falseself._node_sync_event = threading.Event()self._pod_sync_event = threading.Event()def add_or_update_pod(self, pod_dict):with self._lock:uid = pod_dict.get('metadata', {}).get('uid')if uid:self._pods[uid] = pod_dictdef delete_pod(self, pod_dict):with self._lock:uid = pod_dict.get('metadata', {}).get('uid')if uid and uid in self._pods:del self._pods[uid]def get_all_pods(self):with self._lock:return list(self._pods.values())def get_pod_count(self):with self._lock:return len(self._pods)def add_or_update_node(self, node_dict):with self._lock:name = node_dict.get('metadata', {}).get('name')if name:self._nodes[name] = node_dictdef delete_node(self, node_dict):with self._lock:name = node_dict.get('metadata', {}).get('name')if name and name in self._nodes:del self._nodes[name]def get_all_nodes(self):with self._lock:return list(self._nodes.values())def get_node_count(self):with self._lock:return len(self._nodes)def mark_node_synced(self):with self._lock:self._node_synced = Trueself._node_sync_event.set()logger.info("[{}] Node cache synced ({} nodes)".format(self.environment, self.get_node_count()))def mark_pod_synced(self):with self._lock:self._pod_synced = Trueself._pod_sync_event.set()logger.info("[{}] Pod cache synced ({} pods)".format(self.environment, self.get_pod_count()))def wait_for_node_sync(self, timeout=None):return self._node_sync_event.wait(timeout=timeout)def wait_for_pod_sync(self, timeout=None):return self._pod_sync_event.wait(timeout=timeout)def is_node_synced(self):with self._lock:return self._node_synceddef is_pod_synced(self):with self._lock:return self._pod_synceddef clear(self):with self._lock:self._pods.clear()self._nodes.clear()self._node_synced = Falseself._pod_synced = Falseself._node_sync_event.clear()self._pod_sync_event.clear()# ============================================================================ # Reflector (v13 - Raw HTTP 方式) # ============================================================================class Reflector(object):LIST_LIMIT = 5000LIST_TIMEOUT = 60TOTAL_LIST_TIMEOUT = 600def __init__(self, environment, core_api, api_client, resource_type, cache,list_func, watch_stream_func, resync_period=300,on_list_complete=None, namespaces=None, label_selector=None):self.environment = environmentself.core_api = core_apiself.api_client = api_client # v13: 需要 api_client 来调用 raw APIself.resource_type = resource_typeself.cache = cacheself.list_func = list_funcself.watch_stream_func = watch_stream_funcself.resync_period = resync_periodself.on_list_complete = on_list_completeself.namespaces = namespacesself.label_selector = label_selectorself._stop_event = threading.Event()self._thread = Noneself._reconnect_backoff = 1.0self._max_reconnect_backoff = 30.0def start(self):self._thread = threading.Thread(target=self._run,name='reflector-{}-{}'.format(self.environment, self.resource_type),daemon=True)self._thread.start()logger.info("[{}] Reflector for {} started".format(self.environment, self.resource_type))def stop(self):self._stop_event.set()if self._thread and self._thread.is_alive():self._thread.join(timeout=5)def _run(self):while not self._stop_event.is_set():try:if self.resource_type == 'pod' and self.namespaces:success = self._do_list_by_namespaces()else:success = self._do_list_raw()if not success:self._backoff_and_wait()continueif self.on_list_complete:try:self.on_list_complete(self.resource_type)except Exception as e:logger.error("[{}] on_list_complete error: {}".format(self.environment, e))self._do_watch()if self.resync_period 0:logger.debug("[{}] {} watch ended, will resync".format(self.environment, self.resource_type))else:logger.warning("[{}] {} watch disconnected, reconnecting...".format(self.environment, self.resource_type))self._backoff_and_wait()except Exception as e:logger.error("[{}] Reflector {} error: {}".format(self.environment, self.resource_type, e), exc_info=True)self._backoff_and_wait()def _do_list_raw(self):"""v13: 使用 _preload_content=False 获取原始 JSON 响应"""try:logger.info("[{}] Listing all {}s (raw mode)...".format(self.environment, self.resource_type))start_time = time.time()list_deadline = start_time + self.TOTAL_LIST_TIMEOUTall_items = []_continue = Nonepage = 0while True:if self._stop_event.is_set():return Falseif time.time() list_deadline:logger.error("[{}] {} list exceeded total timeout ({}s), aborting".format(self.environment, self.resource_type, self.TOTAL_LIST_TIMEOUT))return Falsepage += 1kwargs = {'limit': self.LIST_LIMIT,'_preload_content': False, # v13: 关键!获取原始 HTTP 响应'_request_timeout': self.LIST_TIMEOUT,}if _continue:kwargs['_continue'] = _continueif self.label_selector:kwargs['label_selector'] = self.label_selectortry:# v13: 调用 API 获取原始响应response = self.list_func(**kwargs)# v13: 解析 JSONdata = json.loads(response.data)except Exception as e:logger.error("[{}] {} list page {} failed: {}".format(self.environment, self.resource_type, page, e))return Falseitems = data.get('items', [])# v13: 直接存储 dict,不创建对象for item in items:if self.resource_type == 'pod':self.cache.add_or_update_pod(item)else:self.cache.add_or_update_node(item)all_items.extend(items)_continue = data.get('metadata', {}).get('continue')if len(all_items) % 1000 len(items) or len(all_items) == len(items):elapsed = time.time() - start_timelogger.info("[{}] {} list progress: {} items in {:.1f}s (page {})".format(self.environment, self.resource_type, len(all_items), elapsed, page))if not _continue:breakelapsed = time.time() - start_timelogger.info("[{}] Listed {} {}s in {:.2f}s ({} pages) [RAW MODE]".format(self.environment, len(all_items), self.resource_type, elapsed, page))self._reconnect_backoff = 1.0return Trueexcept Exception as e:logger.error("[{}] List {} error: {}".format(self.environment, self.resource_type, e))return Falsedef _do_list_by_namespaces(self):"""v13: 按 namespace 分别 list raw"""try:logger.info("[{}] Listing pods by namespaces (raw mode): {}".format(self.environment, self.namespaces))start_time = time.time()list_deadline = start_time + self.TOTAL_LIST_TIMEOUTtotal_items = []for ns in self.namespaces:if self._stop_event.is_set():return Falseif time.time() list_deadline:logger.error("[{}] Pod list exceeded total timeout ({}s), aborting".format(self.environment, self.TOTAL_LIST_TIMEOUT))return Falsens_start = time.time()ns_items = []_continue = Nonepage = 0while True:page += 1kwargs = {'limit': self.LIST_LIMIT,'namespace': ns,'_preload_content': False, # v13: raw mode'_request_timeout': self.LIST_TIMEOUT,}if _continue:kwargs['_continue'] = _continueif self.label_selector:kwargs['label_selector'] = self.label_selectortry:response = self.core_api.list_namespaced_pod(**kwargs)data = json.loads(response.data)except Exception as e:logger.error("[{}] Pod list namespace={} page {} failed: {}".format(self.environment, ns, page, e))return Falseitems = data.get('items', [])for item in items:self.cache.add_or_update_pod(item)ns_items.extend(items)_continue = data.get('metadata', {}).get('continue')if not _continue:breakns_elapsed = time.time() - ns_startlogger.info("[{}] Listed {} pods in namespace={} in {:.2f}s ({} pages) [RAW]".format(self.environment, len(ns_items), ns, ns_elapsed, page))total_items.extend(ns_items)elapsed = time.time() - start_timelogger.info("[{}] Listed {} pods in {} namespaces in {:.2f}s [RAW MODE]".format(self.environment, len(total_items), len(self.namespaces), elapsed))self._reconnect_backoff = 1.0return Trueexcept Exception as e:logger.error("[{}] Pod list by namespaces error: {}".format(self.environment, e))return Falsedef _do_watch(self):try:logger.info("[{}] Watching {} (from latest)".format(self.environment, self.resource_type))w = watch.Watch()kwargs = {}if self.resync_period:kwargs['timeout_seconds'] = self.resync_periodif self.label_selector:kwargs['label_selector'] = self.label_selectorevent_count = 0for event in w.stream(self.watch_stream_func, **kwargs):if self._stop_event.is_set():breakevent_type = event['type']obj = event['object']# v13: namespace 过滤if self.resource_type == 'pod' and self.namespaces:if obj.metadata.namespace not in self.namespaces:continue# v13: 将对象转换为 dict 存储obj_dict = self.api_client.sanitize_for_serialization(obj)if event_type == 'ADDED' or event_type == 'MODIFIED':if self.resource_type == 'pod':self.cache.add_or_update_pod(obj_dict)else:self.cache.add_or_update_node(obj_dict)elif event_type == 'DELETED':if self.resource_type == 'pod':self.cache.delete_pod(obj_dict)else:self.cache.delete_node(obj_dict)event_count += 1if event_count % 100 == 0:logger.debug("[{}] {} processed {} watch events".format(self.environment, self.resource_type, event_count))logger.info("[{}] {} watch ended after {} events".format(self.environment, self.resource_type, event_count))except ApiException as e:if e.status == 410:logger.warning("[{}] {} watch 410 Gone, will re-list".format(self.environment, self.resource_type))else:logger.error("[{}] {} watch API error: status={}, reason={}".format(self.environment, self.resource_type, e.status, e.reason))except Exception as e:logger.error("[{}] {} watch error: {}".format(self.environment, self.resource_type, e))def _backoff_and_wait(self):logger.info("[{}] {} reflector backing off for {:.1f}s".format(self.environment, self.resource_type, self._reconnect_backoff))self._stop_event.wait(self._reconnect_backoff)self._reconnect_backoff = min(self._reconnect_backoff * 2, self._max_reconnect_backoff)# ============================================================================ # Metrics Calculator (v13 - 直接操作 dict) # ============================================================================_POD_STATUSES = ('Pending', 'Running', 'Succeeded', 'Failed', 'Unknown', 'Terminating', 'ContainerCreating')class MetricsCalculator(object):"""v13: 直接操作 dict,不创建对象"""@staticmethoddef _determine_pod_status(pod_dict):metadata = pod_dict.get('metadata', {})if metadata.get('deletionTimestamp'):return 'Terminating'container_statuses = pod_dict.get('status', {}).get('containerStatuses', [])for cs in container_statuses:waiting = cs.get('state', {}).get('waiting')if waiting and waiting.get('reason') == 'ContainerCreating':return 'ContainerCreating'init_statuses = pod_dict.get('status', {}).get('initContainerStatuses', [])for cs in init_statuses:waiting = cs.get('state', {}).get('waiting')if waiting and waiting.get('reason') == 'ContainerCreating':return 'ContainerCreating'phase = pod_dict.get('status', {}).get('phase', 'Unknown')if phase in _POD_STATUSES:return phasereturn 'Unknown'@classmethoddef calculate_and_format(cls, environment, cache):start_time = time.time()lines = []try:# --- Sync Status Metrics ---node_synced = 1 if cache.is_node_synced() else 0pod_synced = 1 if cache.is_pod_synced() else 0lines.append('# HELP scankube_sync_status 1 if the resource type is synced, 0 otherwise')lines.append('# TYPE scankube_sync_status gauge')lines.append('scankube_sync_status{{environment="{}",resource_type="node"}} {}'.format(environment, node_synced))lines.append('scankube_sync_status{{environment="{}",resource_type="pod"}} {}'.format(environment, pod_synced))# --- Node Metrics ---nodes = cache.get_all_nodes()total_nodes = len(nodes)not_ready = 0for node in nodes:conditions = node.get('status', {}).get('conditions', [])is_ready = any(c.get('type') == 'Ready' and c.get('status') == 'True'for c in conditions)if not is_ready:not_ready += 1lines.append('# HELP scankube_node_not_ready Number of NotReady nodes in the cluster')lines.append('# TYPE scankube_node_not_ready gauge')lines.append('scankube_node_not_ready{{environment="{}"}} {}'.format(environment, not_ready))lines.append('# HELP scankube_node_total Total number of nodes in the cluster')lines.append('# TYPE scankube_node_total gauge')lines.append('scankube_node_total{{environment="{}"}} {}'.format(environment, total_nodes))# --- Pod Status Metrics ---pods = cache.get_all_pods()status_counts = defaultdict(lambda: defaultdict(int))for pod in pods:namespace = pod.get('metadata', {}).get('namespace', 'unknown')status = cls._determine_pod_status(pod)status_counts[namespace][status] += 1if status_counts:lines.append('# HELP scankube_pod_status_by_namespace Number of pods by status in each namespace')lines.append('# TYPE scankube_pod_status_by_namespace gauge')for namespace, counts in sorted(status_counts.items()):for status in _POD_STATUSES:count = counts.get(status, 0)lines.append('scankube_pod_status_by_namespace{{environment="{}",namespace="{}",status="{}"}} {}'.format(environment, namespace, status, count))# --- Kaspersky Metrics ---kasp_stats = defaultdict(lambda: {'total': 0, 'terminating': 0, 'container_creating': 0})for pod in pods:pod_name = pod.get('metadata', {}).get('name', '')if 'kaspersky' not in pod_name.lower():continuenamespace = pod.get('metadata', {}).get('namespace', 'unknown')kasp_stats[namespace]['total'] += 1if pod.get('metadata', {}).get('deletionTimestamp'):kasp_stats[namespace]['terminating'] += 1has_cc = Falsefor cs in pod.get('status', {}).get('containerStatuses', []):waiting = cs.get('state', {}).get('waiting')if waiting and waiting.get('reason') == 'ContainerCreating':has_cc = Truebreakif not has_cc:for cs in pod.get('status', {}).get('initContainerStatuses', []):waiting = cs.get('state', {}).get('waiting')if waiting and waiting.get('reason') == 'ContainerCreating':has_cc = Truebreakif has_cc:kasp_stats[namespace]['container_creating'] += 1if kasp_stats:lines.append('# HELP scankube_kaspersky_terminating_pods Number of kaspersky pods in Terminating state')lines.append('# TYPE scankube_kaspersky_terminating_pods gauge')lines.append('# HELP scankube_kaspersky_container_creating_pods Number of kaspersky pods in ContainerCreating state')lines.append('# TYPE scankube_kaspersky_container_creating_pods gauge')lines.append('# HELP scankube_kaspersky_total_pods Total number of kaspersky pods in the namespace')lines.append('# TYPE scankube_kaspersky_total_pods gauge')for namespace, stats in sorted(kasp_stats.items()):lines.append('scankube_kaspersky_terminating_pods{{environment="{}",namespace="{}"}} {}'.format(environment, namespace, stats['terminating']))lines.append('scankube_kaspersky_container_creating_pods{{environment="{}",namespace="{}"}} {}'.format(environment, namespace, stats['container_creating']))lines.append('scankube_kaspersky_total_pods{{environment="{}",namespace="{}"}} {}'.format(environment, namespace, stats['total']))# --- Cache Size ---lines.append('# HELP scankube_cache_size Number of objects in local cache')lines.append('# TYPE scankube_cache_size gauge')lines.append('scankube_cache_size{{environment="{}",resource_type="node"}} {}'.format(environment, total_nodes))lines.append('scankube_cache_size{{environment="{}",resource_type="pod"}} {}'.format(environment, len(pods)))# --- Scrape Metrics ---elapsed = time.time() - start_timelines.append('# HELP scankube_last_scrape_timestamp Unix timestamp of the last scrape')lines.append('# TYPE scankube_last_scrape_timestamp gauge')lines.append('scankube_last_scrape_timestamp{{environment="{}"}} {}'.format(environment, time.time()))lines.append('# HELP scankube_scrape_error 1 if the last scrape had an error, 0 otherwise')lines.append('# TYPE scankube_scrape_error gauge')lines.append('scankube_scrape_error{{environment="{}"}} 0'.format(environment))lines.append('# HELP scankube_scrape_duration_seconds Duration of the last scrape cycle in seconds')lines.append('# TYPE scankube_scrape_duration_seconds gauge')lines.append('scankube_scrape_duration_seconds{{environment="{}"}} {}'.format(environment, round(elapsed, 3)))except Exception as e:logger.error("[{}] Error calculating metrics: {}".format(environment, e), exc_info=True)elapsed = time.time() - start_timelines.append('# HELP scankube_scrape_error 1 if the last scrape had an error, 0 otherwise')lines.append('# TYPE scankube_scrape_error gauge')lines.append('scankube_scrape_error{{environment="{}"}} 1'.format(environment))lines.append('# HELP scankube_scrape_duration_seconds Duration of the last scrape cycle in seconds')lines.append('# TYPE scankube_scrape_duration_seconds gauge')lines.append('scankube_scrape_duration_seconds{{environment="{}"}} {}'.format(environment, round(elapsed, 3)))return '\n'.join(lines) + '\n'# ============================================================================ # Environment Manager - v13 # ============================================================================class EnvironmentManager(object):def __init__(self, kubeconfig_path, namespaces=None, label_selector=None):self.kubeconfig_path = kubeconfig_pathself.environment = os.path.basename(kubeconfig_path)for ext in ['.conf', '.yaml', '.yml', '.kubeconfig']:if self.environment.endswith(ext):self.environment = self.environment[:-len(ext)]breakself.api_client = Noneself.core_api = Noneself.cache = LocalCache(self.environment)self.node_reflector = Noneself.pod_reflector = Noneself._stop_event = threading.Event()self._connected = Falseself._initial_calc_done = Falseself._calc_lock = threading.Lock()self.namespaces = namespacesself.label_selector = label_selectordef connect(self):try:self.api_client = config.new_client_from_config(config_file=self.kubeconfig_path)if self.api_client and self.api_client.rest_client:pool_manager = self.api_client.rest_client.pool_managerif pool_manager:pool_manager.connection_pool_kw['maxsize'] = 20pool_manager.connection_pool_kw['block'] = Falselogger.debug("[{}] Connection pool optimized (maxsize=20)".format(self.environment))self.core_api = client.CoreV1Api(api_client=self.api_client)logger.info("[{}] Connected to k8s cluster".format(self.environment))self._connected = Truereturn Trueexcept Exception as e:logger.error("[{}] Failed to connect: {}".format(self.environment, e))return Falsedef _on_list_complete(self, resource_type):if resource_type == 'node':self.cache.mark_node_synced()logger.debug("[{}] Node synced, {} nodes".format(self.environment, self.cache.get_node_count()))self._trigger_calc_if_needed()elif resource_type == 'pod':self.cache.mark_pod_synced()logger.debug("[{}] Pod synced, {} pods".format(self.environment, self.cache.get_pod_count()))logger.info("[{}] Pod data synced, triggering metrics calc".format(self.environment))controller.trigger_calc(self.environment, force=True)def _trigger_calc_if_needed(self):with self._calc_lock:if self._initial_calc_done:returnself._initial_calc_done = Truelogger.info("[{}] Initial sync completed ({} nodes, {} pods), triggering calc".format(self.environment, self.cache.get_node_count(), self.cache.get_pod_count()))controller.trigger_calc(self.environment)def start(self):if not self._connected:logger.error("[{}] Not connected, cannot start".format(self.environment))return Falseself.node_reflector = Reflector(environment=self.environment,core_api=self.core_api,api_client=self.api_client, # v13resource_type='node',cache=self.cache,list_func=self.core_api.list_node,watch_stream_func=self.core_api.list_node,resync_period=300,on_list_complete=lambda rt: self._on_list_complete(rt))self.node_reflector.start()if self.namespaces:list_func = self.core_api.list_namespaced_podwatch_func = self.core_api.list_namespaced_podelse:list_func = self.core_api.list_pod_for_all_namespaceswatch_func = self.core_api.list_pod_for_all_namespacesself.pod_reflector = Reflector(environment=self.environment,core_api=self.core_api,api_client=self.api_client, # v13resource_type='pod',cache=self.cache,list_func=list_func,watch_stream_func=watch_func,resync_period=300,on_list_complete=lambda rt: self._on_list_complete(rt),namespaces=self.namespaces,label_selector=self.label_selector)self.pod_reflector.start()logger.info("[{}] Environment manager started".format(self.environment))return Truedef stop(self):self._stop_event.set()if self.node_reflector:self.node_reflector.stop()if self.pod_reflector:self.pod_reflector.stop()if self.api_client:try:self.api_client.close()except Exception:passlogger.info("[{}] Environment manager stopped".format(self.environment))def get_cache_stats(self):return {'environment': self.environment,'pods': self.cache.get_pod_count(),'nodes': self.cache.get_node_count(),'node_synced': self.cache.is_node_synced(),'pod_synced': self.cache.is_pod_synced(),}# ============================================================================ # Flask Routes # ============================================================================@app.route('/metrics') def metrics():output = metrics_cache.get_all_output()return Response(output, mimetype='text/plain; version=0.0.4; charset=utf-8')@app.route('/') @app.route('/health') def health():return Response('scankube_monitor_v13 is running\n', mimetype='text/plain')@app.route('/status') def status():stats = []for env, manager in _env_managers.items():stats.append(manager.get_cache_stats())return jsonify({'version': 'v13', 'environments': stats, 'cache': metrics_cache.get_stats()})@app.route('/debug') def debug():debug_info = {}for env, manager in _env_managers.items():debug_info[env] = {'cache': manager.get_cache_stats()}return jsonify(debug_info)# ============================================================================ # Global State # ============================================================================_env_managers = {} _server_started = False _server_lock = threading.Lock()def _start_server(listen_port):global _server_startedwith _server_lock:if _server_started:return_server_started = Trueflask_thread = threading.Thread(target=app.run,kwargs={'host': '0.0.0.0','port': listen_port,'threaded': True,'use_reloader': False,},daemon=True,name='flask-server')flask_thread.start()logger.info("Flask metrics server started on port {}".format(listen_port))# ============================================================================ # Main Entry Point # ============================================================================def main(kubeconfig_path, listen_port=9101, namespaces=None, label_selector=None):manager = EnvironmentManager(kubeconfig_path, namespaces=namespaces, label_selector=label_selector)if not manager.connect():logger.error("Failed to connect: {}".format(kubeconfig_path))return Noneif not manager.start():logger.error("Failed to start: {}".format(kubeconfig_path))return None_env_managers[manager.environment] = managerlogger.info("Added environment: {}".format(manager.environment))_start_server(listen_port)return managerdef serve():if not _env_managers:logger.error("No environments registered. Call main() first.")returncontroller.start()logger.info("scankube_monitor_v13 is running. Press Ctrl+C to stop.")try:while True:time.sleep(1)except KeyboardInterrupt:logger.info("Shutting down...")controller.stop()for manager in list(_env_managers.values()):manager.stop()if __name__ == '__main__':import argparseparser = argparse.ArgumentParser(description='scankube_monitor_v13 - Raw HTTP K8s Monitor')parser.add_argument('--kubeconfig', '-k', action='append', required=True,help='Path to kubeconfig file (can be specified multiple times)')parser.add_argument('--port', '-p', type=int, default=9101,help='HTTP listen port (default: 9101)')parser.add_argument('--workers', '-w', type=int, default=10,help='ThreadPoolExecutor max workers (default: 10)')parser.add_argument('--fallback-interval', '-f', type=int, default=30,help='Fallback calc interval in seconds (default: 30)')parser.add_argument('--namespaces', '-n', type=str, default=None,help='Comma-separated list of namespaces to monitor (default: all)')parser.add_argument('--label-selector', '-l', type=str, default=None,help='Label selector for pods (e.g. app=kaspersky)')args = parser.parse_args()controller._executor = ThreadPoolExecutor(max_workers=args.workers, thread_name_prefix='metrics-calc')controller._fallback_interval = args.fallback_intervalnamespaces = args.namespaces.split(',') if args.namespaces else Nonefor kc_path in args.kubeconfig:main(kubeconfig_path=kc_path, listen_port=args.port, namespaces=namespaces, label_selector=args.label_selector)if not _env_managers:logger.error("No valid connections. Exiting.")sys.exit(1)serve()
← 返回列表