[Bench] fix: Update stress_cluster_benchmark (#791)

Co-authored-by: 春明 <chunming.dc@alibaba-inc.com>
This commit is contained in:
Chaos Dong 2025-08-29 14:46:42 +08:00 committed by GitHub
parent fc04a31cae
commit d944902104
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
1 changed files with 13 additions and 15 deletions

View File

@ -13,6 +13,7 @@ import threading
import queue
import copy
import math
import sys
from collections import defaultdict
# Disable memcpy optimization
@ -78,8 +79,8 @@ class PerformanceTracker:
self.operation_latencies: List[float] = []
self.operation_sizes: List[int] = []
self.error_codes: Dict[int, int] = defaultdict(int)
self.start_time: float = 0
self.end_time: float = 0
self.start_time: float = sys.float_info.max
self.end_time: float = sys.float_info.min
self.total_operations: int = 0
self.failed_operations: int = 0
self.bytes_transferred: int = 0
@ -101,6 +102,8 @@ class PerformanceTracker:
self.total_operations += other.total_operations
self.failed_operations += other.failed_operations
self.bytes_transferred += other.bytes_transferred
self.start_time = min(self.start_time, other.start_time)
self.end_time = max(self.end_time, other.end_time)
# Merge error codes
for code, count in other.error_codes.items():
@ -241,7 +244,7 @@ class TestInstance:
# Prepare batch data
keys = [f"key{total_operations + i}" for i in range(current_batch_size)]
# Measure batch operation latency
# Measure batch operation latency
op_start = time.perf_counter()
return_codes = operation_func(keys, current_batch_size)
op_end = time.perf_counter()
@ -527,7 +530,6 @@ def main():
# Multi-threaded execution
results_queue = queue.Queue()
threads = []
start_time = time.perf_counter()
# Adjust requests per worker
requests_per_worker = args.max_requests // args.num_workers
@ -537,7 +539,8 @@ def main():
for i in range(args.num_workers):
worker_args = copy.copy(args)
worker_args.max_requests = requests_per_worker + (1 if i < remainder else 0)
worker_args.thread_id = i + 1
thread = threading.Thread(
target=worker_thread,
args=(worker_args, results_queue, start_barrier, end_barrier),
@ -550,14 +553,11 @@ def main():
logger.info("Main thread waiting at start barrier")
start_barrier.wait()
logger.info("Main thread passed start barrier")
start_time = time.perf_counter()
# Wait for all workers to complete at end barrier
logger.info("Main thread waiting at end barrier")
end_barrier.wait()
logger.info("Main thread passed end barrier")
end_time = time.perf_counter()
wall_time = end_time - start_time
# Combine results
combined_tracker = PerformanceTracker()
@ -567,17 +567,14 @@ def main():
tracker = results_queue.get()
worker_stats.append(tracker)
combined_tracker.extend(tracker)
# Set combined wall time
combined_tracker.start_time = start_time
combined_tracker.end_time = end_time
# Print detailed worker stats if requested
if args.detailed_stats:
for i, tracker in enumerate(worker_stats):
stats = tracker.get_statistics()
print_performance_stats(stats, f"Worker {i+1} {args.role}")
wall_time = combined_tracker.get_total_time()
# Print combined statistics
combined_stats = combined_tracker.get_statistics()
print_performance_stats(combined_stats, "COMBINED PERFORMANCE")
@ -585,6 +582,7 @@ def main():
# Log precise wall time measurement
logger.info(f"Precise wall time measurement: {wall_time:.4f} seconds")
else:
args.thread_id = 1
# Single-threaded execution
tester = TestInstance(args)
tester.setup()
@ -604,4 +602,4 @@ def main():
if __name__ == '__main__':
main()
main()