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
29 changes: 25 additions & 4 deletions swarms/regression.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
#!/usr/bin/env python3
import os
import random
import sys
from testflows.core import *

Expand All @@ -18,13 +20,23 @@
)
from helpers.feature_support import validate_feature_support, setting_supported


def argparser(parser):
"""Custom argparser that adds swarms-specific options."""
argparser_minio(parser)

parser.add_argument(
"--seed",
action="store",
type=int,
default=None,
help="seed for sampling the join and union combinations, default: $GITHUB_RUN_ID in CI, otherwise random (printed in the log)",
)

from swarms.requirements.requirements import *


xfails = {
"/swarms/feature/swarm joins/join clause/join 455 of 816480*": [
(Fail, "https://github.com/Altinity/ClickHouse/issues/1244"),
],
"/swarms/feature/task rescheduling/rescheduling with bucket granularity": [
(
Fail,
Expand Down Expand Up @@ -64,7 +76,7 @@
@Name("swarms")
@FFails(ffails)
@XFails(xfails)
@ArgumentParser(argparser_minio)
@ArgumentParser(argparser)
@Specifications(SRS_044_Swarm_Cluster_Query_Execution)
@CaptureClusterArgs
@CaptureMinioArgs
Expand All @@ -75,6 +87,7 @@ def regression(
stress=None,
with_analyzer=False,
minio_args=None,
seed=None,
):
"""Run tests for Swarm clusters."""
nodes = {
Expand All @@ -92,6 +105,14 @@ def regression(
if stress is not None:
self.context.stress = stress

# One seed per CI workflow run: every job of the run, and a rerun of
# any of them, samples the same join and union combinations.
if seed is None:
seed = os.environ.get("GITHUB_RUN_ID")
if seed is None:
seed = random.SystemRandom().randrange(2**32)
self.context.seed = int(seed)

minio_root_user = minio_args["minio_root_user"].value
minio_root_password = minio_args["minio_root_password"].value

Expand Down
83 changes: 50 additions & 33 deletions swarms/tests/joins.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,9 +27,12 @@ def __init__(
cluster_name=None,
minio_root_user=None,
minio_root_password=None,
label=None,
):
self.location = location
self.table_type = table_type
# Stable name for the scenario name: the table names carry getuid().
self.label = label
self.cluster_name = cluster_name
self.minio_root_user = minio_root_user or JoinTable.minio_root_user
self.minio_root_password = minio_root_password or JoinTable.minio_root_password
Expand Down Expand Up @@ -538,7 +541,7 @@ def join_clause(self, minio_root_user, minio_root_password, node=None):
)

with Given("create iceberg tables in different locations"):
for location in locations:
for i, location in enumerate(locations, 1):
_, table_name, namespace = (
swarm_steps.iceberg_table_with_all_basic_data_types(
minio_root_user=minio_root_user,
Expand All @@ -551,6 +554,7 @@ def join_clause(self, minio_root_user, minio_root_password, node=None):
database_name=database_name,
namespace=namespace,
table_name=table_name,
label=f"iceberg_table_data{i}",
)
)

Expand All @@ -561,37 +565,46 @@ def join_clause(self, minio_root_user, minio_root_password, node=None):
url,
minio_root_user=minio_root_user,
minio_root_password=minio_root_password,
label=f"s3_data{i}",
)
for url in urls
for i, url in enumerate(urls, 1)
]
iceberg_table_functions = [
JoinTable.create_iceberg_table_function(
url,
minio_root_user=minio_root_user,
minio_root_password=minio_root_password,
label=f"iceberg_data{i}",
)
for url in urls
for i, url in enumerate(urls, 1)
]
iceberg_s3_cluster_table_functions = [
JoinTable.create_icebergS3Cluster_table_function(url, "replicated_cluster")
for url in urls
JoinTable.create_icebergS3Cluster_table_function(
url, "replicated_cluster", label=f"icebergS3Cluster_data{i}"
)
for i, url in enumerate(urls, 1)
]
s3_cluster_table_functions = [
JoinTable.create_s3Cluster_table_function(url, "replicated_cluster")
for url in urls
JoinTable.create_s3Cluster_table_function(
url, "replicated_cluster", label=f"s3Cluster_data{i}"
)
for i, url in enumerate(urls, 1)
]

with Given(
"create merge tree tables from iceberg tables with same schema and data"
):
merge_tree_tables = []
for iceberg_table in iceberg_tables:
for i, iceberg_table in enumerate(iceberg_tables, 1):
merge_tree_table_name = f"merge_tree_table_{getuid()}"
create_table_as_select(
as_select_from=iceberg_table, table_name=merge_tree_table_name
)
merge_tree_tables.append(
JoinTable.create_merge_tree_table(table_name=merge_tree_table_name)
JoinTable.create_merge_tree_table(
table_name=merge_tree_table_name,
label=f"merge_tree_table_data{i}",
)
)

left_tables = (
Expand Down Expand Up @@ -652,19 +665,6 @@ def join_clause(self, minio_root_user, minio_root_password, node=None):
"FULL OUTER JOIN",
]

length = len(
list(
product(
left_tables,
right_tables,
modes,
join_conditions,
object_storage_clusters,
join_clauses,
)
)
)

all_possible_combinations = list(
product(
left_tables,
Expand All @@ -676,21 +676,38 @@ def join_clause(self, minio_root_user, minio_root_password, node=None):
)
)

# A scenario is named after its combination: "join N of <total>" is the
# position of the combination in the full product, so "join N" is the
# same combination in every run. Each run executes a different sample of
# combinations, drawn from the run's seed (see regression.py), so the
# numbers of a run are not sequential. To reproduce a failure, use the
# seed printed below, otherwise "join N" is probably not in the sample:
# --seed <seed> --only "/swarms/feature/swarm joins/join clause/join N of*"
total = len(all_possible_combinations)
combination_ids = range(total)

if not self.context.stress:
all_possible_combinations = random.sample(
all_possible_combinations, min(1000, len(all_possible_combinations))
note(f"join combinations seed: {self.context.seed}")
combination_ids = sorted(
random.Random(self.context.seed).sample(combination_ids, min(1000, total))
)

with Pool() as pool:
for num, (
left_table,
right_table,
mode,
join_condition,
object_storage_cluster,
join_clause,
) in enumerate(all_possible_combinations):
name = f"join {num} of {length}: {left_table} with {right_table} in {mode} mode on {object_storage_cluster} cluster with {join_clause} clause"
for combination_id in combination_ids:
(
left_table,
right_table,
mode,
join_condition,
object_storage_cluster,
join_clause,
) = all_possible_combinations[combination_id]
condition = join_condition.replace("t1.", "").replace("t2.", "")
name = (
f"join {combination_id} of {total}: {left_table.label} with {right_table.label}"
f" in {mode} mode on {object_storage_cluster} cluster"
f" with {join_clause} on {condition}"
)
Scenario(
name=name,
test=check_join,
Expand Down
Loading