Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 9 additions & 7 deletions asap-tools/dataset-analysis/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ Two tasks:

```bash
pip install -r requirements.txt
./fetch_data.sh /path/to/trace-data # ~63 GB; re-run to resume
./fetch_data.sh /path/to/trace-data # ~207 GB; re-run to resume
python fit_skew.py --data-root /path/to/trace-data
```

Expand All @@ -37,19 +37,21 @@ This writes `results/skew_summary.csv` (committed) and plots to `out/`
- `--workers N`: process pool size (default: all cores). Files and fits run
in parallel.

A full run over the fetched data takes about 2 hours with 24 workers on a
56-core machine (peak RSS of the main process about 69 GB). Alibaba archives are streamed with `tarfile`, not extracted.
With 48 workers on a 56-core, 251 GB machine, Google takes about 1 hour
(peak RSS of the main process 45 GB) and Alibaba about 10.5 hours (98 GB;
CallGraph's 6h windows over high-cardinality keys dominate). Run the two
datasets on separate machines with `--queries` to overlap them. Alibaba archives are streamed with `tarfile`, not extracted.

Tests: `python -m unittest discover -s tests -p 'test_*.py'`.

## Datasets

| Dataset | Files | Tables | Step | Ranges |
|---|---|---|---|---|
| Google ClusterData 2011-2 | `task_usage` and `task_events` parts 0..119 of 500 (about 7 days), all 500 `job_events` parts, `schema.csv` | `task_usage` joined with task and job attributes | 5 min | instant, 5m, 1h |
| Alibaba microservices v2022 | `MCRRTUpdate_0..119` (first 6 hours) | `MSRTMCR` | 1 min | instant, 5m, 1h |
| | `CallGraph_0..119` (first 6 hours) | `CallGraph` | 1 min | 1m, 5m, 1h (events, no instant) |
| | `MSMetricsUpdate_0..47`, `NodeMetricsUpdate_0..1` (first day) | `MSMetrics`, `NodeMetrics` | 1 min | instant, 5m, 1h |
| Google ClusterData 2011-2 | `task_usage` and `task_events` parts 0..119 of 500 (about 7 days), all 500 `job_events` parts, `schema.csv` | `task_usage` joined with task and job attributes | 5 min | instant, 5m, 1h, 6h, 24h |
| Alibaba microservices v2022 | `MCRRTUpdate_0..479` (first day) | `MSRTMCR` | 1 min | instant, 5m, 1h, 6h, 24h |
| | `CallGraph_0..479` (first day) | `CallGraph` | 1 min | 1m, 5m, 1h, 6h, 24h (events, no instant) |
| | `MSMetricsUpdate_0..47` (first day), `NodeMetricsUpdate_0..2` (first 1.5 days: the first shard starts a step late, so two hold no full 24h window) | `MSMetrics`, `NodeMetrics` | 1 min | instant, 5m, 1h, 6h, 24h |
| Datadog BOOM | `dataset_taxonomy.json` and 20 multivariate series | per-series `target` | none | 20 equal chunks per series |

Citations:
Expand Down
9 changes: 5 additions & 4 deletions asap-tools/dataset-analysis/fetch_data.sh
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,12 @@ BOOM_URL=https://huggingface.co/datasets/Datadog/BOOM/resolve/main
# task_events uses the same parts; job_events is read in full for job names.
GOOGLE_TASK_PARTS=120
GOOGLE_JOB_EVENT_PARTS=500
# CallGraph and MCRRTUpdate shards cover 3 minutes each: 120 shards = 6 hours.
ALIBABA_RPC_SHARDS=120
# MSMetricsUpdate shards cover 30 minutes, NodeMetricsUpdate 12 hours: 1 day.
# CallGraph and MCRRTUpdate shards cover 3 minutes each: 480 shards = 1 day.
ALIBABA_RPC_SHARDS=480
# MSMetricsUpdate shards cover 30 minutes: 48 = 1 day. NodeMetricsUpdate shards
# cover 12 hours, but the first starts a step late: 3 hold a full 24h window.
ALIBABA_MS_SHARDS=48
ALIBABA_NODE_SHARDS=2
ALIBABA_NODE_SHARDS=3

BOOM_SERIES=(
ds-2187-H ds-2394-D ds-1135-5T ds-1833-D ds-2806-D
Expand Down
68 changes: 53 additions & 15 deletions asap-tools/dataset-analysis/fit_skew.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@
ZIPF_THETA_BOUNDS = (0.0, 5.0)
# Each power-law fit uses a uniform subsample of at most this many values.
MAX_FIT_SAMPLES = 100_000
# numpy's multivariate_hypergeometric needs a total under 1e9 (merge_samples).
HYPERGEOMETRIC_MAX_TOTAL = 1_000_000_000
# xmin is chosen by KS distance over this many quantiles of the sample, and only
# where at least MIN_TAIL_SAMPLES values remain in the tail.
XMIN_GRID_SIZE = 50
Expand All @@ -64,6 +66,8 @@
DURATION_UNITS_S = {"s": 1, "m": 60, "h": 3600, "d": 86400}
# Range key sums are materialized for this many evaluation times at a time.
EVAL_CHUNK = 32
# Steps per chunk when stitching instant samples across files (resolve_boundaries).
BOUNDARY_CHUNK_STEPS = 60
# Accuracy each sketch must reach on a query; a query's `targets` overrides.
DEFAULT_TARGETS = {
"are_top100": 0.05, # CMS / CountSketch mean relative error of the top 100 keys
Expand Down Expand Up @@ -432,11 +436,18 @@ def value_sample(x: np.ndarray) -> Tuple[int, np.ndarray]:

def merge_samples(parts: Sequence[Tuple[int, np.ndarray]]) -> Tuple[int, np.ndarray]:
"""Uniform sample of the union of value_sample parts: each part gives a
multivariate-hypergeometric share of its prefix."""
multivariate-hypergeometric share of its prefix. numpy draws that only
for totals under 1e9; above, a multinomial share (drawing 1e5 of 1e9 or
more, with and without replacement agree), capped at each prefix."""
counts = np.array([n for n, _ in parts], dtype=np.int64)
total = int(counts.sum())
rng = np.random.default_rng(SAMPLE_SEED)
take = rng.multivariate_hypergeometric(counts, min(total, MAX_FIT_SAMPLES))
n = min(total, MAX_FIT_SAMPLES)
if total < HYPERGEOMETRIC_MAX_TOTAL:
take = rng.multivariate_hypergeometric(counts, n)
else:
lengths = np.array([len(x) for _, x in parts])
take = np.minimum(rng.multinomial(n, counts / total), lengths)
merged = np.concatenate([x[:k] for (_, x), k in zip(parts, take)])
return total, rng.permutation(merged)

Expand Down Expand Up @@ -526,22 +537,49 @@ def check_file_order(rows: pd.DataFrame) -> None:


def resolve_boundaries(
boundary: Sequence[pd.DataFrame], queries: List[Dict[str, Any]], lookback: int
boundary: Sequence[pd.DataFrame],
queries: List[Dict[str, Any]],
lookback: int,
chunk_steps: int = BOUNDARY_CHUNK_STEPS,
) -> Dict[Tuple[str, str], Any]:
"""Instant aggregates of the per-file first/last samples: keep the latest
sample per (series, step) over files, and take its next step as the
nearer of the in-file next step and the next boundary sample."""
rows = pd.concat(
[b.assign(**{FILE_COL: i}) for i, b in enumerate(boundary)], ignore_index=True
)
check_file_order(rows)
rows = rows.sort_values([SERIES_COL, STEP_COL, TIME_COL], kind="stable")
in_file_next = rows.groupby([SERIES_COL, STEP_COL])[NEXT_COL].transform("min")
rows = rows.assign(**{NEXT_COL: in_file_next})
rows = rows.drop_duplicates([SERIES_COL, STEP_COL], keep="last")
next_step, _ = next_step_of_series(rows)
rows = rows.assign(**{NEXT_COL: np.fmin(rows[NEXT_COL].to_numpy(), next_step)})
return instant_parts(rows, queries, lookback)
nearer of the in-file next step and the next boundary sample.

Short files make almost every sample a boundary one, so this runs over
chunks of `chunk_steps` steps. A sample counts for at most `lookback`
steps, so a chunk owning steps [start, end) reads the boundary samples
up to end + lookback and gives the same aggregates as one pass."""
files = [b.assign(**{FILE_COL: i}) for i, b in enumerate(boundary) if len(b)]
if not files:
return instant_parts(pd.concat(boundary, ignore_index=True), queries, lookback)
spans = [(int(f[STEP_COL].min()), int(f[STEP_COL].max())) for f in files]
first, last = min(lo for lo, _ in spans), max(hi for _, hi in spans)
chunks = []
for start in range(first, last + 1, chunk_steps):
end = start + chunk_steps
overlapping = [
f[(f[STEP_COL] >= start) & (f[STEP_COL] < end + lookback)]
for f, (lo, hi) in zip(files, spans)
if lo < end + lookback and hi >= start
]
if not overlapping:
continue
rows = pd.concat(overlapping, ignore_index=True)
check_file_order(rows)
rows = rows.sort_values([SERIES_COL, STEP_COL, TIME_COL], kind="stable")
in_file_next = rows.groupby([SERIES_COL, STEP_COL])[NEXT_COL].transform("min")
rows = rows.assign(**{NEXT_COL: in_file_next})
rows = rows.drop_duplicates([SERIES_COL, STEP_COL], keep="last")
next_step, _ = next_step_of_series(rows)
rows = rows.assign(**{NEXT_COL: np.fmin(rows[NEXT_COL].to_numpy(), next_step)})
chunks.append(instant_parts(rows[rows[STEP_COL] < end], queries, lookback))
return {
query_key(q): (merge_key_parts if q["kind"] == "keys" else merge_step_samples)(
[c[query_key(q)] for c in chunks]
)
for q in queries
}


def instant_parts(
Expand Down
37 changes: 18 additions & 19 deletions asap-tools/dataset-analysis/queries/alibaba_v2022.yaml
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
# Alibaba cluster-trace-microservices-v2022: first 6 hours of CallGraph and
# MCRRTUpdate, first day of MSMetricsUpdate and NodeMetricsUpdate.
# Alibaba cluster-trace-microservices-v2022: the first day of every table.
#
# Tables: step_s is the sampling period and the evaluation step; series_key
# (metric tables only) is the label set identifying one series.
Expand Down Expand Up @@ -48,7 +47,7 @@ queries:
- id: mcr_by_msname
promql: sum by (msname) (providerrpc_mcr)
promql_range: sum by (msname) (sum_over_time(providerrpc_mcr[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSRTMCR
kind: keys
group_by: [msname]
Expand All @@ -58,7 +57,7 @@ queries:
- id: mcr_by_nodeid
promql: sum by (nodeid) (providerrpc_mcr)
promql_range: sum by (nodeid) (sum_over_time(providerrpc_mcr[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSRTMCR
kind: keys
group_by: [nodeid]
Expand All @@ -68,7 +67,7 @@ queries:
- id: rt_p99_by_msname
promql: quantile by (msname) (0.99, providerrpc_rt)
promql_range: quantile by (msname) (0.99, quantile_over_time(0.99, providerrpc_rt[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSRTMCR
kind: keys
group_by: [msname]
Expand All @@ -77,80 +76,80 @@ queries:
- id: rt_p99_by_msname
promql: quantile(0.99, providerrpc_rt)
promql_range: quantile_over_time(0.99, providerrpc_rt[{range}])
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSRTMCR
kind: values
group_by: []
value: providerrpc_rt
# CallGraph rows are events (one per call), so only range queries apply.
- id: calls_by_rpctype
promql_range: sum by (rpctype) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [rpctype]
weights: [count]
- id: calls_by_service
promql_range: sum by (service) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [service]
weights: [count]
- id: calls_by_um
promql_range: sum by (um) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [um]
weights: [count]
- id: calls_by_dm
promql_range: sum by (dm) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [dm]
weights: [count]
- id: calls_by_interface
promql_range: sum by (interface) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [interface]
weights: [count]
- id: calls_by_um_dm
promql_range: sum by (um, dm) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [um, dm]
weights: [count]
- id: calls_by_service_um_dm
promql_range: sum by (service, um, dm) (count_over_time(rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [service, um, dm]
weights: [count]
- id: rt_p99_by_service
promql_range: quantile by (service) (0.99, quantile_over_time(0.99, rt[{range}]))
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: keys
group_by: [service]
value: rt
weights: [count, value]
- id: rt_p99_by_service
promql_range: quantile_over_time(0.99, rt[{range}])
range: [1m, 5m, 1h]
range: [1m, 5m, 1h, 6h, 24h]
table: CallGraph
kind: values
group_by: []
value: rt
- id: ms_cpu_by_msname
promql: sum by (msname) (cpu_utilization)
promql_range: sum by (msname) (sum_over_time(cpu_utilization[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSMetrics
kind: keys
group_by: [msname]
Expand All @@ -161,7 +160,7 @@ queries:
- id: node_cpu_by_nodeid
promql: sum by (nodeid) (cpu_utilization)
promql_range: sum by (nodeid) (sum_over_time(cpu_utilization[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: NodeMetrics
kind: keys
group_by: [nodeid]
Expand All @@ -171,15 +170,15 @@ queries:
- id: ms_cpu_p99
promql: quantile(0.99, cpu_utilization)
promql_range: quantile_over_time(0.99, cpu_utilization[{range}])
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: MSMetrics
kind: values
group_by: []
value: cpu_utilization
- id: node_cpu_p99
promql: quantile(0.99, cpu_utilization)
promql_range: quantile_over_time(0.99, cpu_utilization[{range}])
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: NodeMetrics
kind: values
group_by: []
Expand Down
18 changes: 9 additions & 9 deletions asap-tools/dataset-analysis/queries/google_2011.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ queries:
- id: cpu_by_user
promql: sum by (user) (cpu_rate)
promql_range: sum by (user) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [user]
Expand All @@ -37,7 +37,7 @@ queries:
- id: cpu_by_user_priority
promql: sum by (user, priority) (cpu_rate)
promql_range: sum by (user, priority) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [user, priority]
Expand All @@ -47,7 +47,7 @@ queries:
- id: cpu_by_priority
promql: sum by (priority) (cpu_rate)
promql_range: sum by (priority) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [priority]
Expand All @@ -57,7 +57,7 @@ queries:
- id: cpu_by_job_id
promql: sum by (job_id) (cpu_rate)
promql_range: sum by (job_id) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [job_id]
Expand All @@ -67,7 +67,7 @@ queries:
- id: cpu_by_logical_job_name
promql: sum by (logical_job_name) (cpu_rate)
promql_range: sum by (logical_job_name) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [logical_job_name]
Expand All @@ -77,7 +77,7 @@ queries:
- id: cpu_by_scheduling_class
promql: sum by (scheduling_class) (cpu_rate)
promql_range: sum by (scheduling_class) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [scheduling_class]
Expand All @@ -88,7 +88,7 @@ queries:
- id: cpu_by_machine_id
promql: sum by (machine_id) (cpu_rate)
promql_range: sum by (machine_id) (sum_over_time(cpu_rate[{range}]))
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: keys
group_by: [machine_id]
Expand All @@ -98,15 +98,15 @@ queries:
- id: cpu_p99
promql: quantile(0.99, cpu_rate)
promql_range: quantile_over_time(0.99, cpu_rate[{range}])
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: values
group_by: []
value: cpu_rate
- id: memory_p99
promql: quantile(0.99, canonical_memory_usage)
promql_range: quantile_over_time(0.99, canonical_memory_usage[{range}])
range: [instant, 5m, 1h]
range: [instant, 5m, 1h, 6h, 24h]
table: task_usage
kind: values
group_by: []
Expand Down
Loading
Loading