"""InfluxDB 备份:通过 influx / influxd CLI 生成备份目录,再 tar.gz 流式打包。 支持 v1.x 和 v2.x: - v1.x:使用 influx backup (CLI 子命令) 等价做法:用 influx_inspect export 转 line protocol,然后打包 这里采用更通用的:influxd backup(如果有) 或 打包 line-protocol dump - v2.x:使用 influx backup 命令生成备份目录 策略: 1. 在临时目录执行 backup 命令生成文件 2. 用 tarfile 流式打包整个目录 3. 子进程产出后整体打包上传 """ import asyncio import os import shutil import subprocess import tarfile import tempfile from datetime import datetime from pathlib import Path from typing import AsyncIterator from app.core.backup.base import BackupMetadata, BaseBackup from app.utils.logging import get_logger log = get_logger(__name__) class InfluxDBBackup(BaseBackup): def __init__(self, source_config: dict, job_name: str): super().__init__(source_config, job_name) self.version: str = source_config["version"] self.host: str = source_config["host"] self.port: int = int(source_config.get("port", 8086)) if self.version == "1": self.database: str = source_config["database"] self.username: str = source_config.get("username", "") self.password: str = source_config.get("password", "") else: self.org: str = source_config["org"] self.bucket: str = source_config["bucket"] self.token: str = source_config["token"] def metadata(self) -> BackupMetadata: ts = datetime.now().strftime("%Y%m%d-%H%M%S") if self.version == "1": safe = "".join(c if c.isalnum() or c in ("_", "-") else "_" for c in self.database) else: safe = "".join(c if c.isalnum() or c in ("_", "-") else "_" for c in self.bucket) return BackupMetadata( suggested_filename=f"{self.job_name}-{self.version}-{safe}-{ts}.tar.gz", content_type="application/gzip", extra={"version": self.version, "host": self.host, "port": self.port}, ) async def produce(self) -> AsyncIterator[bytes]: ts = datetime.now().strftime("%Y%m%d%H%M%S") tmpdir = Path(tempfile.mkdtemp(prefix=f"influxbk_{self.job_name}_")) backup_dir = tmpdir / "data" backup_dir.mkdir() try: if self.version == "2": await self._backup_v2(backup_dir) else: await self._backup_v1(backup_dir) # 打包为 tar.gz 并流式输出 log.info("Taring influx backup at %s", backup_dir) tar_path = tmpdir / "backup.tar.gz" def _tar() -> None: with tarfile.open(str(tar_path), mode="w:gz") as tar: tar.add(str(backup_dir), arcname="data", recursive=True) loop = asyncio.get_event_loop() await loop.run_in_executor(None, _tar) if not tar_path.exists() or tar_path.stat().st_size == 0: raise RuntimeError("InfluxDB 备份为空,请检查数据库连接和 CLI 安装") # 读取打包文件并流式产出 with open(tar_path, "rb") as f: while True: chunk = await loop.run_in_executor(None, f.read, 64 * 1024) if not chunk: break yield chunk finally: shutil.rmtree(tmpdir, ignore_errors=True) async def _backup_v2(self, dest: Path) -> None: """v2: influx backup --bucket X --org Y /dest""" cmd = [ "influx", "backup", "--bucket", self.bucket, "--org", self.org, "--host", f"http://{self.host}:{self.port}", "--token", self.token, str(dest), ] log.info("Running: %s", " ".join(self._redact(cmd))) proc = await asyncio.create_subprocess_exec( *cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, ) stdout, stderr = await proc.communicate() if proc.returncode != 0: raise RuntimeError( f"influx backup failed (exit={proc.returncode}): " f"{(stderr or stdout).decode(errors='replace')[:2000]}" ) async def _backup_v1(self, dest: Path) -> None: """v1: 通过 HTTP API 查询全量数据,导出为 CSV。 使用 httpx 直接调 InfluxDB v1 的 /query 接口(更可靠,不依赖 CLI 子命令兼容性)。 生成: - {database}.csv:SELECT * 结果 - META.txt:恢复所需的元信息(库名、host、时间戳等) 恢复:需手动用 influx -import 或写入 line protocol。 """ import httpx log.info("InfluxDB v1 backup via HTTP API: db=%s host=%s:%s", self.database, self.host, self.port) url = f"http://{self.host}:{self.port}/query" params = { "db": self.database, "q": "SELECT * FROM /.*/", "chunked": "true", "epoch": "ns", } auth = None if self.username: auth = (self.username, self.password or "") async with httpx.AsyncClient(timeout=300.0, auth=auth) as client: try: resp = await client.get(url, params=params) except Exception as e: raise RuntimeError(f"InfluxDB v1 HTTP 请求失败: {e}") from e if resp.status_code != 200: raise RuntimeError( f"InfluxDB v1 返回 {resp.status_code}: {resp.text[:500]}" ) body = resp.text # 保存 CSV 和元信息 csv_file = dest / f"{self.database}.csv" loop = asyncio.get_event_loop() def _write_files() -> None: csv_file.write_text(body, encoding="utf-8") meta = dest / "META.txt" meta.write_text( f"influxdb_version=1\n" f"host={self.host}\nport={self.port}\n" f"database={self.database}\n" f"username={self.username or ''}\n" f"query=SELECT * FROM /.*/\n" f"format=csv\n" f"backup_time={datetime.utcnow().isoformat()}Z\n" f"restore_hint=可使用 `influx -import -path={csv_file.name} " f"-database={self.database}` 恢复;或在 InfluxQL shell 中执行 " f"`INSERT INTO ... VALUES ...`\n", encoding="utf-8", ) await loop.run_in_executor(None, _write_files) if not csv_file.exists() or csv_file.stat().st_size == 0: raise RuntimeError("InfluxDB v1 查询结果为空,请检查数据库是否存在或是否有数据") @staticmethod def _redact(cmd: list[str]) -> list[str]: """日志脱敏 token/password。""" out = [] skip_next = False secret_flags = {"--token", "-password"} for i, c in enumerate(cmd): if skip_next: out.append("***") skip_next = False continue if c in secret_flags: out.append(c) skip_next = True continue out.append(c) return out