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
10 changes: 10 additions & 0 deletions tests/integration/test_backup_restore_new/configs/cas_gc.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
<clickhouse>
<storage_configuration>
<disks>
<cas>
<cas_gc_enabled>1</cas_gc_enabled>
<cas_gc_interval_sec>1</cas_gc_interval_sec>
</cas>
</disks>
</storage_configuration>
</clickhouse>
24 changes: 24 additions & 0 deletions tests/integration/test_backup_restore_new/configs/cas_storage.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
<clickhouse>
<storage_configuration>
<disks>
<cas>
<type>object_storage</type>
<object_storage_type>s3</object_storage_type>
<metadata_type>cas</metadata_type>
<cas_server_root_id>itest-backup-restore-new</cas_server_root_id>
<endpoint>http://rustfs1:11121/test/backup_restore_new/</endpoint>
<access_key_id>clickhouse</access_key_id>
<secret_access_key>clickhouse</secret_access_key>
</cas>
</disks>
<policies>
<cas_policy>
<volumes>
<single>
<disk>cas</disk>
</single>
</volumes>
</cas_policy>
</policies>
</storage_configuration>
</clickhouse>
99 changes: 90 additions & 9 deletions tests/integration/test_backup_restore_new/test.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,10 @@
cluster = ClickHouseCluster(__file__)
instance = cluster.add_instance(
"instance",
main_configs=["configs/backups_disk.xml"],
main_configs=["configs/backups_disk.xml", "configs/cas_storage.xml"],
user_configs=["configs/zookeeper_retries.xml"],
external_dirs=["/backups/"],
with_rustfs=True,
)
instance_with_short_timeout = cluster.add_instance(
"instance_with_short_timeout",
Expand All @@ -30,11 +31,14 @@
)


def create_and_fill_table(engine="MergeTree", n=100):
def create_and_fill_table(engine="MergeTree", n=100, storage_policy=None):
if engine == "MergeTree":
engine = "MergeTree ORDER BY y PARTITION BY x%10"
instance.query("CREATE DATABASE test")
instance.query(f"CREATE TABLE test.table(x UInt32, y String) ENGINE={engine}")
create_query = f"CREATE TABLE test.table(x UInt32, y String) ENGINE={engine}"
if storage_policy is not None:
create_query += f" SETTINGS storage_policy = '{storage_policy}'"
instance.query(create_query)
instance.query(
f"INSERT INTO test.table SELECT number, toString(number) FROM numbers({n})"
)
Expand Down Expand Up @@ -224,6 +228,30 @@ def test_restore_table(engine):
assert instance.query("SELECT count(), sum(x) FROM test.table") == "100\t4950\n"


def test_restore_table_on_cas_disk():
backup_name = new_backup_name()
create_and_fill_table(storage_policy="cas_policy")

assert instance.query("SELECT count(), sum(x) FROM test.table") == "100\t4950\n"
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== "cas_policy\n"
)
instance.query(f"BACKUP TABLE test.table TO {backup_name}")
instance.query("DROP TABLE test.table")
instance.query(f"RESTORE TABLE test.table FROM {backup_name}")

assert instance.query("SELECT count(), sum(x) FROM test.table") == "100\t4950\n"
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== "cas_policy\n"
)
assert (
instance.query("CHECK TABLE test.table SETTINGS check_query_single_value_result = 1")
== "1\n"
)


@pytest.mark.parametrize(
"engine", ["MergeTree", "Log", "TinyLog", "StripeLog", "Memory"]
)
Expand Down Expand Up @@ -1504,8 +1532,9 @@ def test_system_users_async():
)


def test_projection():
create_and_fill_table(n=3)
@pytest.mark.parametrize("storage_policy", [None, "cas_policy"])
def test_projection(storage_policy):
create_and_fill_table(n=3, storage_policy=storage_policy)

instance.query("ALTER TABLE test.table ADD PROJECTION prjmax (SELECT MAX(x))")
instance.query("INSERT INTO test.table VALUES (100, 'a'), (101, 'b')")
Expand Down Expand Up @@ -1555,6 +1584,18 @@ def test_projection():
== "2\n"
)

if storage_policy is not None:
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== f"{storage_policy}\n"
)
assert (
instance.query(
"CHECK TABLE test.table SETTINGS check_query_single_value_result = 1"
)
== "1\n"
)


def test_restore_table_not_evaluate_table_defaults():
instance.query("CREATE DATABASE test")
Expand Down Expand Up @@ -1950,8 +1991,10 @@ def verify_restore_info():
assert info.bytes_read == 0


def test_mutation():
create_and_fill_table(engine="MergeTree ORDER BY tuple()", n=5)
def check_mutation(storage_policy=None):
create_and_fill_table(
engine="MergeTree ORDER BY tuple()", n=5, storage_policy=storage_policy
)

instance.query(
"INSERT INTO test.table SELECT number, toString(number) FROM numbers(5, 5)"
Expand All @@ -1973,6 +2016,12 @@ def test_mutation():
instance.query("ALTER TABLE test.table UPDATE x=x+1 WHERE 1")
instance.query("ALTER TABLE test.table UPDATE x=x+1 WHERE 1")

if storage_policy is not None:
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== f"{storage_policy}\n"
)
assert instance.query("SELECT count() FROM test.table") == "15\n"
backup_name = new_backup_name()
instance.query(f"BACKUP TABLE test.table TO {backup_name}")

Expand All @@ -1986,6 +2035,20 @@ def test_mutation():
instance.query("DROP TABLE test.table")

instance.query(f"RESTORE TABLE test.table FROM {backup_name}")
if storage_policy is not None:
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== f"{storage_policy}\n"
)
assert instance.query("SELECT count() FROM test.table") == "15\n"


def test_mutation():
check_mutation()


def test_mutation_on_cas_disk():
check_mutation(storage_policy="cas_policy")


def test_tables_dependency():
Expand Down Expand Up @@ -2215,12 +2278,17 @@ def test_restore_table_with_checksum_data_file_name(engine):
assert instance.query("SELECT count(), sum(x) FROM test.table") == "100\t4950\n"


def test_incremental_backup_with_checksum_data_file_name():
def check_incremental_backup_with_checksum_data_file_name(storage_policy=None):
backup_name = new_backup_name()
incremental_backup_name = new_backup_name()
create_and_fill_table()
create_and_fill_table(storage_policy=storage_policy)

assert instance.query("SELECT count(), sum(x) FROM test.table") == "100\t4950\n"
if storage_policy is not None:
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table'")
== f"{storage_policy}\n"
)
instance.query(
f"BACKUP TABLE test.table TO {backup_name} SETTINGS data_file_name_generator='checksum'"
)
Expand All @@ -2236,6 +2304,19 @@ def test_incremental_backup_with_checksum_data_file_name():
f"RESTORE TABLE test.table AS test.table2 FROM {incremental_backup_name}"
)
assert instance.query("SELECT count(), sum(x) FROM test.table2") == "102\t5081\n"
if storage_policy is not None:
assert (
instance.query("SELECT storage_policy FROM system.tables WHERE name = 'table2'")
== f"{storage_policy}\n"
)


def test_incremental_backup_with_checksum_data_file_name():
check_incremental_backup_with_checksum_data_file_name()


def test_incremental_backup_with_checksum_data_file_name_on_cas_disk():
check_incremental_backup_with_checksum_data_file_name(storage_policy="cas_policy")


def test_async_backup_restore_with_max_execution_time_zero():
Expand Down
73 changes: 71 additions & 2 deletions tests/integration/test_backup_restore_new/test_cancel_backup.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import re
import time
import uuid

import pytest
Expand All @@ -17,15 +18,20 @@
"configs/backups_disk.xml",
"configs/slow_backups.xml",
"configs/shutdown_cancel_backups.xml",
"configs/cas_storage.xml",
"configs/cas_gc.xml",
]

node = cluster.add_instance(
"node",
main_configs=main_configs,
external_dirs=["/backups/"],
stay_alive=True,
with_rustfs=True,
)

CAS_POOL_PREFIXES = ["backup_restore_new/blobs/"]


@pytest.fixture(scope="module", autouse=True)
def start_cluster():
Expand Down Expand Up @@ -179,11 +185,57 @@ def cancel_restore(restore_id):
assert kill_duration_ms < kill_duration_ms_threshold


def table_settings(storage_policy):
if storage_policy is None:
return ""
return f" SETTINGS storage_policy = '{storage_policy}'"


def count_cas_pool_objects():
return sum(
len(
list(
cluster.rustfs_client.list_objects(
cluster.rustfs_bucket, prefix, recursive=True
)
)
)
for prefix in CAS_POOL_PREFIXES
)


def wait_cas_pool_settled():
previous = count_cas_pool_objects()
for _ in range(40):
time.sleep(3)
current = count_cas_pool_objects()
if current == previous:
return current
previous = current
return previous


def wait_cas_pool_reclaimed(baseline):
current = count_cas_pool_objects()
for _ in range(120):
if current <= baseline:
break
time.sleep(1)
current = count_cas_pool_objects()
assert (
current <= baseline
), f"CAS pool objects were not reclaimed: baseline={baseline}, current={current}"


# Test that BACKUP and RESTORE operations can be cancelled with KILL QUERY.
def test_cancel_backup():
@pytest.mark.parametrize("storage_policy", [None, "cas_policy"])
def test_cancel_backup(storage_policy):
cas_pool_baseline = wait_cas_pool_settled() if storage_policy == "cas_policy" else None

# We use partitioning so backups would contain more files.
node.query(
"CREATE TABLE tbl (x UInt64) ENGINE=MergeTree() ORDER BY tuple() PARTITION BY x%20"
+ table_settings(storage_policy)
)

node.query("INSERT INTO tbl SELECT number FROM numbers(500)")
Expand All @@ -209,11 +261,25 @@ def test_cancel_backup():
start_restore(restore_id, backup_id)
wait_restore(restore_id)

assert node.query("SELECT count(), sum(x) FROM tbl") == "500\t124750\n"
assert (
node.query("CHECK TABLE tbl SETTINGS check_query_single_value_result = 1")
== "1\n"
)

if cas_pool_baseline is not None:
node.query("DROP TABLE tbl SYNC")
wait_cas_pool_reclaimed(cas_pool_baseline)


# Test that shutdown cancels a running backup and doesn't wait until it finishes.
def test_shutdown_cancel_backup():
@pytest.mark.parametrize("storage_policy", [None, "cas_policy"])
def test_shutdown_cancel_backup(storage_policy):
cas_pool_baseline = wait_cas_pool_settled() if storage_policy == "cas_policy" else None

node.query(
"CREATE TABLE tbl (x UInt64) ENGINE=MergeTree() ORDER BY tuple() PARTITION BY x%5"
+ table_settings(storage_policy)
)

node.query("INSERT INTO tbl SELECT number FROM numbers(500)")
Expand All @@ -237,3 +303,6 @@ def test_shutdown_cancel_backup():
f"RESTORE TABLE tbl FROM {get_backup_name(backup_id)}"
),
)

if cas_pool_baseline is not None:
wait_cas_pool_reclaimed(cas_pool_baseline)
Original file line number Diff line number Diff line change
Expand Up @@ -38,17 +38,50 @@ def generate_cluster_def(file: str, num_nodes: int) -> str:
return str(path.absolute())


def generate_cas_server_root_def(file: str, node_index: int) -> str:
path = (
Path(__file__).parent
/ f"_gen/cas_server_root_{Path(file).stem}_node{node_index}.xml"
)
config = f"""<clickhouse>
<storage_configuration>
<disks>
<cas>
<cas_server_root_id>itest-{Path(file).stem}-node{node_index}</cas_server_root_id>
</cas>
</disks>
</storage_configuration>
</clickhouse>"""
if path.is_file() and path.read_text(encoding="utf-8") == config:
return str(path.absolute())
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(encoding="utf-8", data=config)
return str(path.absolute())


def add_nodes_to_cluster(
cluster: ClickHouseCluster,
num_nodes: int,
main_configs: List[str],
user_configs: List[str],
cas_file: str = None,
**kwargs
) -> List[ClickHouseInstance]:
def node_main_configs(i):
if cas_file is None:
return main_configs
return main_configs + [
"configs/cas_storage.xml",
generate_cas_server_root_def(cas_file, i),
]

if cas_file is not None:
kwargs["with_rustfs"] = True

nodes = [
cluster.add_instance(
f"node{i}",
main_configs=main_configs,
main_configs=node_main_configs(i),
user_configs=user_configs,
external_dirs=["/backups/"],
macros={"replica": f"node{i}", "shard": "shard1"},
Expand All @@ -60,9 +93,11 @@ def add_nodes_to_cluster(
return nodes


def create_test_table(node: ClickHouseInstance) -> None:
def create_test_table(node: ClickHouseInstance, storage_policy: str = None) -> None:
settings = f" SETTINGS storage_policy = '{storage_policy}'" if storage_policy else ""
node.query(
"""CREATE TABLE tbl ON CLUSTER 'cluster' ( x UInt64 )
ENGINE=ReplicatedMergeTree('/clickhouse/tables/tbl/', '{replica}')
ORDER BY tuple()"""
+ settings
)
Loading
Loading