-
Notifications
You must be signed in to change notification settings - Fork 1
AUTO_BUFFER
Stand: 15. Dezember 2025
Version: 1.0.0
Status: Design/Implementierung
Kategorie: Technical Documentation
Der TSAutoBuffer ist ein konfigurierbarer Buffering-Layer für TSStore, der automatisch einzelne Datenpunkte sammelt und als komprimierte Batches speichert. Dies löst die Limitierung, dass putDataPoint() keine Gorilla-Kompression nutzt.
Aktueller Zustand:
-
putDataPoint(): Speichert Punkte als einzelne RocksDB-Entities (keine Kompression) -
putDataPoints(): Speichert Batches mit Gorilla-Kompression (10-20x Ratio) - HTTP-API
/ts/putnutztputDataPoint()→ keine Kompression
Nachteil:
- IoT/Streaming-Anwendungen können nicht von Gorilla-Kompression profitieren
- Einzelpunkt-Einfügungen verschwenden Speicherplatz
Ein intelligenter Buffer der:
- Einzelne Datenpunkte puffert
- Automatisch nach Größe/Zeit flusht
- Gorilla-Kompression nutzt
- Thread-safe ist
- Mit bestehender TSStore-API arbeitet
┌─────────────────────────────────────────────────────┐
│ Application │
└────────────────────┬────────────────────────────────┘
│
│ add(point)
▼
┌─────────────────────────────────────────────────────┐
│ TSAutoBuffer │
│ ┌─────────────────────────────────────────────┐ │
│ │ Buffers (map<metric:entity, deque<Point>>) │ │
│ │ - Per-metric:entity buffering │ │
│ │ - Configurable thresholds │ │
│ │ - Thread-safe access │ │
│ └─────────────────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────┐ │
│ │ Background Flush Thread │ │
│ │ - Time-based flush (every N seconds) │ │
│ │ - Size-based flush (max points reached) │ │
│ │ - Memory-based flush (max memory reached) │ │
│ └─────────────────────────────────────────────┘ │
└────────────────────┬────────────────────────────────┘
│
│ putDataPoints(batch)
▼
┌─────────────────────────────────────────────────────┐
│ TSStore │
│ - Gorilla compression │
│ - RocksDB storage │
└─────────────────────────────────────────────────────┘
Der Buffer flusht automatisch wenn:
-
Size-Threshold:
max_points_per_buffererreicht (z.B. 1000 Punkte) -
Time-Threshold:
flush_intervalabgelaufen (z.B. 5 Sekunden) -
Memory-Threshold:
max_memory_byteserreicht (z.B. 100 MB) -
Global-Threshold:
max_total_pointsüber alle Buffers erreicht
struct TSAutoBufferConfig {
// Buffer-Größe
size_t max_points_per_buffer = 1000; // Max Punkte pro metric:entity
size_t max_total_points = 10000; // Max Punkte gesamt
// Zeit-basiert
std::chrono::milliseconds flush_interval{5000}; // 5 Sekunden
// Speicher
size_t max_memory_bytes = 100 * 1024 * 1024; // 100 MB
// Performance
bool async_flush = true; // Background-Thread
size_t flush_batch_size = 500; // Punkte pro Flush
// Kompression (von TSStore)
TSStore::CompressionType compression = TSStore::CompressionType::Gorilla;
int chunk_size_hours = 24;
};TSAutoBufferConfig config;
config.max_points_per_buffer = 500; // Klein für niedrige Latenz
config.flush_interval = std::chrono::seconds(2); // Schnelles Flush
config.max_memory_bytes = 50 * 1024 * 1024; // 50 MB
config.async_flush = true;
config.compression = TSStore::CompressionType::Gorilla;Eigenschaften:
- Niedrige Latenz (2s Flush)
- Kleine Batches für schnelles Schreiben
- Kompression aktiv
TSAutoBufferConfig config;
config.max_points_per_buffer = 5000; // Große Batches
config.flush_interval = std::chrono::seconds(30); // Seltenes Flush
config.max_memory_bytes = 500 * 1024 * 1024; // 500 MB
config.async_flush = true;
config.compression = TSStore::CompressionType::Gorilla;Eigenschaften:
- Hoher Durchsatz
- Große Batches für maximale Kompression
- Viel Speicher für Pufferung
TSAutoBufferConfig config;
config.max_points_per_buffer = 200; // Klein
config.flush_interval = std::chrono::seconds(5);
config.max_memory_bytes = 10 * 1024 * 1024; // Nur 10 MB
config.async_flush = true;
config.compression = TSStore::CompressionType::Gorilla;Eigenschaften:
- Wenig Speicher
- Häufiges Flush
- Trotzdem Kompression
#include "timeseries/tsstore.h"
#include "timeseries/ts_auto_buffer.h"
// TSStore einrichten
TSStore::Config ts_config;
ts_config.compression = TSStore::CompressionType::Gorilla;
TSStore ts(db, cf, ts_config);
// Auto-Buffer einrichten
TSAutoBufferConfig buffer_config;
buffer_config.max_points_per_buffer = 1000;
buffer_config.flush_interval = std::chrono::seconds(5);
TSAutoBuffer buffer(&ts, buffer_config);
buffer.start(); // Background-Thread starten
// Punkte hinzufügen (werden automatisch gepuffert)
for (int i = 0; i < 10000; ++i) {
TSStore::DataPoint point{
.metric = "temperature",
.entity = "sensor_001",
.timestamp_ms = now() + i * 1000,
.value = 20.0 + (rand() % 100) / 10.0,
.tags = {{"location", "room_a"}},
.metadata = {}
};
buffer.add(point); // Gepuffert, nicht direkt geschrieben
}
// Stoppen (flusht automatisch verbleibende Punkte)
buffer.stop();// Alle Buffer flushen
size_t flushed = buffer.flush();
std::cout << "Flushed " << flushed << " points\n";
// Spezifischen Buffer flushen
size_t flushed = buffer.flushFor("temperature", "sensor_001");auto stats = buffer.getStats();
std::cout << "Points buffered: " << stats.points_buffered << "\n";
std::cout << "Points flushed: " << stats.points_flushed << "\n";
std::cout << "Flush count: " << stats.flush_count << "\n";
std::cout << "Auto flushes: " << stats.auto_flush_count << "\n";
std::cout << "Current buffer size: " << stats.current_buffer_size << "\n";
std::cout << "Current memory: " << stats.current_buffer_memory << " bytes\n";// Konfiguration zur Laufzeit ändern
TSAutoBufferConfig new_config;
new_config.max_points_per_buffer = 2000;
new_config.flush_interval = std::chrono::seconds(10);
buffer.setConfig(new_config);Der HTTP-Server kann einen globalen TSAutoBuffer verwenden:
// In http_server.cpp
class HTTPServer {
private:
TSStore tsstore_;
TSAutoBuffer ts_buffer_; // Global buffer
public:
HTTPServer() : ts_buffer_(&tsstore_) {
ts_buffer_.start();
}
void handleTSPut(const Request& req, Response& res) {
auto point = parseDataPoint(req.body());
// Verwende Buffer statt direktem TSStore
auto status = ts_buffer_.add(point);
if (status.ok) {
res.status = 201;
res.body = R"({"success": true})";
} else {
res.status = 400;
res.body = R"({"error": ")" + status.message + R"("})";
}
}
};void handleTSPut(const Request& req, Response& res) {
auto point = parseDataPoint(req.body());
// Check header: X-TS-Buffer: enable
if (req.headers["X-TS-Buffer"] == "enable") {
ts_buffer_.add(point); // Buffered
} else {
tsstore_.putDataPoint(point); // Direct
}
}// POST /ts/put - Direct (kein Buffer)
void handleTSPut(const Request& req, Response& res) {
tsstore_.putDataPoint(parseDataPoint(req.body()));
}
// POST /ts/put/buffered - Mit Auto-Buffer
void handleTSPutBuffered(const Request& req, Response& res) {
ts_buffer_.add(parseDataPoint(req.body()));
}| Modus | Schreiblatenz | Flush-Latenz | Gesamt |
|---|---|---|---|
Direct (putDataPoint) |
~1ms | 0ms | ~1ms |
| Buffered (kleine Batches) | <0.1ms | ~10-50ms | ~0.1ms (async) |
| Buffered (große Batches) | <0.1ms | ~50-200ms | ~0.1ms (async) |
Vorteil: Schreiben ist ~10x schneller (nur Buffer-Operation)
| Modus | Punkte/s | Kompression | Speicher |
|---|---|---|---|
| Direct | ~1,000 | Keine | 100% |
| Buffered (500 pts/batch) | ~10,000 | 10-20x | 5-10% |
| Buffered (5000 pts/batch) | ~50,000 | 15-25x | 4-7% |
Vorteil: Durchsatz ~10-50x höher, Speicher ~90-96% weniger
Buffer Memory = num_buffers × avg_points × point_size
≈ 100 metrics × 1000 points × 200 bytes
≈ 20 MB
Konfigurierbar via max_memory_bytes
Der TSAutoBuffer ist vollständig thread-safe:
- Alle öffentlichen Methoden sind mutex-geschützt
- Background-Thread koordiniert via
std::condition_variable - Atomare Zähler für Statistiken
- Safe Shutdown mit Flush aller Punkte
Beispiel (Multi-Threaded):
TSAutoBuffer buffer(&ts);
buffer.start();
// Thread 1
std::thread t1([&buffer]() {
for (int i = 0; i < 1000; ++i) {
buffer.add(makePoint("cpu", "server1", i));
}
});
// Thread 2
std::thread t2([&buffer]() {
for (int i = 0; i < 1000; ++i) {
buffer.add(makePoint("memory", "server1", i));
}
});
t1.join();
t2.join();
buffer.stop(); // Flusht alle verbleibenden PunkteDie Implementierung nutzt bestehende ThemisDB-Patterns:
-
std::deque<Event> bufferfür Events - Partitionierung nach Event-Typ
- Wiederverwendet: Deque-basierte Pufferung
- Load-basierte Flush-Strategien (IMMEDIATE, THROTTLED, DEFERRED)
- Adaptive Thresholds
- Wiederverwendet: Multi-Threshold Flush-Logik
- Bereits in
putDataPoints()verwendet - Atomare Batch-Schreiboperationen
-
Wiederverwendet: Direkt via
tsstore_->putDataPoints()
-
std::thread,std::mutex,std::condition_variable -
std::atomicfür lockfreie Zähler - Standard-Bibliothek: Keine zusätzlichen Dependencies
Pro:
- Weniger Code
- Direkte RocksDB-Integration
Contra:
- Kein automatisches Flush
- Keine Zeit-basierte Logik
- Keine Statistiken
Pro:
- Lock-free Performance
Contra:
- Zusätzliche Dependency (Boost)
- Komplexere Integration
- Overkill für diesen Use-Case
Pro:
- ✅ Nutzt bestehende Patterns (CEP, Backpressure)
- ✅ Keine neuen Dependencies
- ✅ Volle Kontrolle über Flush-Strategien
- ✅ Integrierte Statistiken
- ✅ Thread-safe Design
- ✅ Einfache API
Contra:
- Etwas mehr Code als Option A
- Speicher-Overhead für Buffers
Entscheidung: Option C ist die beste Balance aus Features, Performance und Wartbarkeit.
- Header-Datei (
ts_auto_buffer.h) - Implementation (
ts_auto_buffer.cpp) - Dokumentation (diese Datei)
- Unit-Tests (
test_ts_auto_buffer.cpp) - Performance-Tests
- Multi-Threading-Tests
- Memory-Leak-Tests
- HTTP-API Integration
- Config-File Support
- Metrics/Tracing
- Documentation Updates
- Lock-free Statistiken
- NUMA-aware Partitioning
- Adaptive Flush-Intervalle
- Predictive Buffering
Frage: Gibt es etwas in den Libs für single-point-to-batch?
Antwort:
- Nein, es gibt keine fertige Bibliothek dafür
- Aber: Bestehende Patterns (CEP, Backpressure) liefern alle Bausteine
-
Lösung:
TSAutoBuffer- konfigurierbarer Auto-Batching Buffer -
Bibliotheken verwendet:
- Standard C++ (
std::deque,std::thread,std::mutex) - RocksDB (
WriteBatchviaTSStore::putDataPoints) - ThemisDB Utils (Logging, Tracing)
- Standard C++ (
- Neue Dependencies: Keine!
Status: Design fertig, Implementation bereit für Testing und Integration.
- Architecture-ACCESS-MODEL-IMPLEMENTATION-SUMMARY
- Architecture-ADR-003-pg-dump-sql-parser
- Architecture-BASEENTITY-PRINCIPLE
- Architecture-CACHE-STORAGE-INTEGRATION
- Architecture-CMAKE-ARCHITECTURE
- Architecture-CMAKE-FLAGS-REFERENCE
- Architecture-CMAKE-MODULAR-ARCHITECTURE
- Architecture-CONCERNS-ARCHITECTURE-DIAGRAM
- Architecture-CONCERNS-IMPLEMENTATION-SUMMARY
- Architecture-CONTENT-MODEL
- Architecture-COPILOT-THEMISDB-GRAPH-RAG-BACKEND-ARCHITECTURE
- Architecture-CRYPTO-AND-KEYS
- Architecture-FEATURE-FLAGS-REFERENCE
- Architecture-GPU-ARCHITECTURE-REVIEW-TEMPLATE
- Architecture-HTTP-SHUTDOWN-HARDENING
- Architecture-MIGRATION-GUIDE-CONCERNS
- Architecture-MIGRATION-GUIDE-v13-v14
- Architecture-MODULARIZATION-GUIDE
- Architecture-MODULAR-ARCHITECTURE-ROADMAP
- Architecture-MODULE-ARCHITECTURE-INDEX
- Architecture-P1D01-ISSMPLUGIN-DESIGN-REVIEW
- Architecture-P1-D01-ISSMPLUGIN-DESIGN-REVIEW
- Architecture-P1-D08-MAMBA-GOVERNANCE-CONTRACT
- Architecture-P1-P2-IMPLEMENTATION-COMPLETION-INDEX
- Architecture-PHASE0-COMPLETION-ASSESSMENT
- Architecture-PHASE3-QUERYENGINE-DI-ARCHITECTURE
- Architecture-PHASE4-INDEX-MANAGER-DI
- Architecture-POSTGRESQL-WIRE-PROTOCOL
- Architecture-QUERYENGINE-IMPLEMENTATION-GUIDE
- Architecture-QUERY-SCHEDULING
- Architecture-RAFT-CONSENSUS-DESIGN
- Architecture-README
- Architecture-README-SSM-HYBRID-IMPLEMENTATION
- Architecture-REFACTORING-SUMMARY
- Architecture-RESOURCE-POOLING
- Architecture-SOURCE-DIRECTORY-GUIDE
- Architecture-THEMIS-CORE-GUIDE
- Architecture-UNIFIED-ACCESS-MODEL
- Architecture-WAL-GRPC-MTLS-CONFIGURATION
- Architecture-WIRE-PROTOCOL-RETRY
- Architecture-boltzmann-observability-draft
- Architecture-experimental-logarithmic-vector-storage
- Architecture-llm-wiki-mvp-adr
- Architecture-rewrite-engine-architecture
- Architecture-rope-api-architecture
- Architecture-ssm-gguf-mamba-status
- Architecture-ssm-hybrid-analysis
- Architecture-ssm-hybrid-rollout-plan
- Architecture-ssm-plugin-interface-design-review
- Architecture-transaction-coordinators
- Architecture-wiki-secondary-index
- Architecture-wire-protocol
- Governance-DISABLED-STUB-POLICY
- Governance-DOCS-PR-POLICY
- Governance-GA-PROMOTION-SIGN-OFF
- Governance-GITHUB-MILESTONES-SETUP
- Governance-MATURITY-CLAIM-VERIFICATION-CHECKLIST
- Governance-MATURITY-EVIDENCE-REGISTRY
- Governance-MERGE-GATE-BOT-CONFIG
- Governance-MERGE-GATE-STATUS-LIVE
- Governance-PHASE3-ENFORCEMENT-RUNBOOK
- Governance-PHASE-1-CLOSURE-REPORT
- Governance-PHASE-CLOSURE-POLICY
- Governance-PHASE-DEPENDENCY-GRAPH
- Governance-PLUGIN-SUBMODULE-ROLLBACK
- Governance-PRODUCTION-READY-2026-DELIVERY-PLAN
- Governance-PR-VERSION-TARGETING
- Governance-PR-VERSION-TARGETING-BACKFILL
- Governance-QUERY-MODULE-STATUS
- Governance-README
- Governance-RELEASE-PROMOTION-GATE-POLICY
- Governance-RELEASE-VALIDATION-CHECKLIST
- Governance-SECURITY-MODULE-5671-EVIDENCE-SUMMARY
- Governance-SHARDING-P6-RESIDUAL-RISK-ACCEPTANCE
- Governance-SOURCECODE-COMPLIANCE-GOVERNANCE
- Governance-UPDATES-DEVELOPMENT-STATUS-SIGN-OFF
- Governance-WAVE-C-IMPLEMENTATION-COMPLETE
- Module-acceleration-Roadmap
- Module-access-model-Roadmap
- Module-ai-Roadmap
- Module-analytics-Roadmap
- Module-api-Roadmap
- Module-aql-Roadmap
- Module-auth-Roadmap
- Module-base-Roadmap
- Module-cache-Roadmap
- Module-cdc-Roadmap
- Module-chaos-Roadmap
- Module-chimera-Roadmap
- Module-config-Roadmap
- Module-content-Roadmap
- Module-core-Roadmap
- Module-distributed-knowledge-Roadmap
- Module-distributed-tensor-Roadmap
- Module-document-Roadmap
- Module-ethics-ai-Roadmap
- Module-evaluation-Roadmap
- Module-execution-Roadmap
- Module-exporters-Roadmap
- Module-failover-Roadmap
- Module-geo-Roadmap
- Module-governance-Roadmap
- Module-gpu-Roadmap
- Module-graph-Roadmap
- Module-image-analysis-Roadmap
- Module-importers-Roadmap
- Module-index-Roadmap
- Module-ingestion-Roadmap
- Module-llama-cpp-Roadmap
- Module-llm-Roadmap
- Module-llm-streaming-Roadmap
- Module-llm-wiki-Roadmap
- Module-maintenance-Roadmap
- Module-metadata-Roadmap
- Module-network-Roadmap
- Module-observability-Roadmap
- Module-onnx-clip-Roadmap
- Module-performance-Roadmap
- Module-plugins-Roadmap
- Module-process-Roadmap
- Module-projects-Roadmap
- Module-prompt-engineering-Roadmap
- Module-query-Roadmap
- Module-rag-Roadmap
- Module-replication-Roadmap
- Module-retrieval-Roadmap
- Module-rpc-grpc-Roadmap
- Module-scheduler-Roadmap
- Module-scraper-Roadmap
- Module-search-Roadmap
- Module-security-Roadmap
- Module-server-Roadmap
- Module-sharding-Roadmap
- Module-stable-diffusion-Roadmap
- Module-storage-Roadmap
- Module-temporal-Roadmap
- Module-tensor-Roadmap
- Module-themis-Roadmap
- Module-timeseries-Roadmap
- Module-toolbox-Roadmap
- Module-training-Roadmap
- Module-transaction-Roadmap
- Module-updates-Roadmap
- Module-user-storage-encrypted-Roadmap
- Module-utils-Roadmap
- Module-vector-search-Roadmap
- Module-voice-Roadmap
- Module-whisper-Roadmap