Skip to content

Commit 82fa620

Browse files
Merge branch 'antalya-26.6' into feature/antalya-26.6/CAS
2 parents cc9034e + fac5bbf commit 82fa620

169 files changed

Lines changed: 2191 additions & 764 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎.github/actions/create_workflow_report/create_workflow_report.py‎

Lines changed: 94 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -691,6 +691,10 @@ def get_new_fails_this_pr(
691691

692692
# Combine both types of fails and select only desired columns
693693
desired_columns = ["job_name", "test_name", "test_status", "results_link"]
694+
if len(checks_fails) > 0 and "labels" in checks_fails.columns:
695+
desired_columns.insert(desired_columns.index("results_link"), "labels")
696+
if len(regression_fails) > 0:
697+
regression_fails["labels"] = ""
694698
all_pr_fails = pd.concat([checks_fails, regression_fails], ignore_index=True)[
695699
desired_columns
696700
]
@@ -976,6 +980,82 @@ def format_test_status(text: str) -> str:
976980
return f'<span style="font-weight: bold; color: {color}">{text}</span>'
977981

978982

983+
def _label_names_from_ext(ext: dict) -> list[str]:
984+
names = []
985+
for item in ext.get("labels") or []:
986+
if isinstance(item, str):
987+
name = item
988+
elif isinstance(item, dict) and item.get("name"):
989+
name = item["name"]
990+
else:
991+
continue
992+
if name != "cidb":
993+
names.append(name)
994+
return names
995+
996+
997+
def fetch_workflow_result_json(
998+
pr_number: int, branch: str, commit_sha: str
999+
) -> dict | None:
1000+
if pr_number == 0:
1001+
ref_param = f"REF={branch}"
1002+
workflow_name = "MasterCI"
1003+
else:
1004+
ref_param = f"PR={pr_number}"
1005+
workflow_name = "PR"
1006+
1007+
status_file = f"result_{workflow_name.lower()}.json"
1008+
s3_path = (
1009+
f"https://{S3_BUCKET}.s3.amazonaws.com/"
1010+
f"{ref_param.replace('=', 's/')}/{commit_sha}/{status_file}"
1011+
)
1012+
try:
1013+
response = requests.get(s3_path, timeout=30)
1014+
if response.status_code != 200:
1015+
return None
1016+
return response.json()
1017+
except Exception as e:
1018+
print(f"WARNING:Failed to fetch workflow result from {s3_path}: {e}")
1019+
return None
1020+
1021+
1022+
def get_failure_labels_from_workflow(workflow_data: dict | None) -> dict:
1023+
if not workflow_data:
1024+
return {}
1025+
labels_map = {}
1026+
for job in workflow_data.get("results") or []:
1027+
job_name = job.get("name")
1028+
if not job_name:
1029+
continue
1030+
for leaf in job.get("results") or []:
1031+
test_name = leaf.get("name")
1032+
if not test_name:
1033+
continue
1034+
names = _label_names_from_ext(leaf.get("ext") or {})
1035+
if names:
1036+
labels_map[(job_name, test_name)] = ", ".join(names)
1037+
return labels_map
1038+
1039+
1040+
def add_labels_to_checks_fails(
1041+
checks_fails: pd.DataFrame, workflow_data: dict | None
1042+
) -> pd.DataFrame:
1043+
if checks_fails is None or len(checks_fails) == 0:
1044+
return checks_fails
1045+
labels_map = get_failure_labels_from_workflow(workflow_data)
1046+
df = checks_fails.copy()
1047+
df["labels"] = df.apply(
1048+
lambda row: labels_map.get((row["job_name"], row["test_name"]), ""),
1049+
axis=1,
1050+
)
1051+
cols = [c for c in df.columns if c != "labels"]
1052+
if "results_link" in cols:
1053+
cols.insert(cols.index("results_link"), "labels")
1054+
else:
1055+
cols.append("labels")
1056+
return df[cols]
1057+
1058+
9791059
def format_results_as_html_table(results, *, branch_name: str = "") -> str:
9801060
if not isinstance(results, pd.DataFrame):
9811061
return results
@@ -1012,42 +1092,30 @@ def format_col_name(col_name: str) -> str:
10121092
"PR Labels": lambda labels: format_pr_labels_with_verification(
10131093
labels, branch_name=branch_name
10141094
),
1095+
"Labels": lambda labels: html.escape(str(labels), quote=True) if labels else "",
10151096
}
10161097

1017-
html = results.to_html(
1098+
return results.to_html(
10181099
index=False,
10191100
formatters=formatters,
10201101
escape=False,
10211102
border=0,
10221103
classes=["test-results-table"],
10231104
)
1024-
return html
10251105

10261106

10271107
def backfill_skipped_statuses(
1028-
job_statuses: pd.DataFrame, pr_number: int, branch: str, commit_sha: str
1108+
job_statuses: pd.DataFrame,
1109+
workflow_result: dict | None,
10291110
):
10301111
"""
10311112
Fill in the job statuses for skipped jobs.
10321113
"""
1033-
1034-
if pr_number == 0:
1035-
ref_param = f"REF={branch}"
1036-
workflow_name = "MasterCI"
1037-
else:
1038-
ref_param = f"PR={pr_number}"
1039-
workflow_name = "PR"
1040-
1041-
status_file = f"result_{workflow_name.lower()}.json"
1042-
s3_path = f"https://{S3_BUCKET}.s3.amazonaws.com/{ref_param.replace('=', 's/')}/{commit_sha}/{status_file}"
1043-
response = requests.get(s3_path)
1044-
1045-
if response.status_code != 200:
1114+
if workflow_result is None:
10461115
return job_statuses
10471116

1048-
status_data = response.json()
10491117
skipped_jobs = []
1050-
for job in status_data["results"]:
1118+
for job in workflow_result["results"]:
10511119
if job["status"] == "skipped" and len(job["links"]) > 0:
10521120
skipped_jobs.append(
10531121
{
@@ -1192,10 +1260,15 @@ def create_workflow_report(
11921260
settings={"use_numpy": True},
11931261
)
11941262

1263+
workflow_result = fetch_workflow_result_json(pr_number, branch_name, commit_sha)
1264+
11951265
results_dfs = {
11961266
"prs_in_release": [],
11971267
"job_statuses": get_commit_statuses(commit_sha),
1198-
"checks_fails": get_checks_fails(db_client, commit_sha, branch_name),
1268+
"checks_fails": add_labels_to_checks_fails(
1269+
get_checks_fails(db_client, commit_sha, branch_name),
1270+
workflow_result,
1271+
),
11991272
"checks_known_fails": [],
12001273
"pr_new_fails": [],
12011274
"checks_errors": get_checks_errors(db_client, commit_sha, branch_name),
@@ -1257,7 +1330,8 @@ def create_workflow_report(
12571330
pr_info = {}
12581331

12591332
results_dfs["job_statuses"] = backfill_skipped_statuses(
1260-
results_dfs["job_statuses"], pr_number, branch_name, commit_sha
1333+
results_dfs["job_statuses"],
1334+
workflow_result,
12611335
)
12621336

12631337
high_cve_count = 0

‎.github/actions/docker_setup/action.yml‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,8 @@ runs:
2020
sudo cat <<EOT > /etc/docker/daemon.json
2121
{
2222
"ipv6": true,
23-
"fixed-cidr-v6": "${{ env.ipv6_subnet }}"
23+
"fixed-cidr-v6": "${{ env.ipv6_subnet }}",
24+
"dns": ["1.1.1.1", "1.0.0.1"]
2425
}
2526
EOT
2627
sudo chown root:root /etc/docker/daemon.json

‎.github/actions/runner_setup/action.yml‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,13 @@ description: Setup environment
33
runs:
44
using: "composite"
55
steps:
6+
- name: Configure docker proxy and dataset mirror host
7+
shell: bash
8+
run: |
9+
proxy_host=dockerhub-proxy.dockerhub-proxy-zone
10+
if ! timeout 5 getent ahostsv4 "$proxy_host" >/dev/null; then
11+
echo "65.108.242.32 $proxy_host" | sudo tee -a /etc/hosts
12+
fi
613
- name: Setup zram
714
shell: bash
815
run: |

‎ci/jobs/stress_job.py‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import csv
22
import logging
33
import os
4+
import socket
45
import sys
56
from pathlib import Path
67
from typing import List, Tuple
@@ -169,12 +170,24 @@ def get_run_command(
169170
else:
170171
run_script = "/repo/tests/docker_scripts/stress_runner.sh"
171172

173+
# Nested docker does not inherit the runner /etc/hosts.
174+
# --network=host would expose minio and azurite. Skip --add-host when the
175+
# name is unknown (CI Tests e2e / local) rather than failing the job.
176+
proxy_host = "dockerhub-proxy.dockerhub-proxy-zone"
177+
add_host = ""
178+
try:
179+
proxy_ip = socket.getaddrinfo(proxy_host, None)[0][4][0]
180+
add_host = f"--add-host={proxy_host}:{proxy_ip} "
181+
except OSError:
182+
logging.info("Could not resolve %s; not passing --add-host", proxy_host)
183+
172184
cmd = (
173185
"docker run --cap-add=SYS_PTRACE "
174186
# For dmesg and sysctl
175187
"--privileged "
176188
# azurite-rs (in-process Azure Blob Storage emulator) needs many fds under parallel load
177189
"--ulimit nofile=1048576:1048576 "
190+
f"{add_host}"
178191
# a static link, don't use S3_URL or S3_DOWNLOAD
179192
"-e S3_URL='https://s3.amazonaws.com/clickhouse-datasets' "
180193
"--tmpfs /tmp/clickhouse:mode=1777 "

‎ci/praktika/runner.py‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import os
66
import re
77
import shlex
8+
import socket
89
import sys
910
import traceback
1011
from pathlib import Path
@@ -446,6 +447,18 @@ def _run(
446447
)
447448
from_root = "root" in docker_settings
448449
settings = [s for s in docker_settings if s.startswith("-")]
450+
# NOTE (strtgbb): FT/integration docker is not --network=host, so
451+
# the runner /etc/hosts mapping for the web-disk proxy is invisible
452+
# inside the container. Resolve on the host and pass --add-host.
453+
# Skip if the name is not configured on this runner (e.g. CI Tests).
454+
_proxy_host = "dockerhub-proxy.dockerhub-proxy-zone"
455+
try:
456+
_proxy_ip = socket.getaddrinfo(_proxy_host, None)[0][4][0]
457+
settings.append(f"--add-host={_proxy_host}:{_proxy_ip}")
458+
except OSError as e:
459+
print(
460+
f"NOTE: skipping --add-host for {_proxy_host}: {e}"
461+
)
449462
if ":" in job.run_in_docker:
450463
docker_name, docker_tag = job.run_in_docker.split(":")
451464
print(

‎docs/en/antalya/part_export.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -53,10 +53,10 @@ Source and destination tables must support positional schema conversion. The fol
5353
- **Column types** may differ, as long as the source type is safely castable to the destination type. Set `export_merge_tree_part_allow_lossy_cast = 1` to also permit lossy casts.
5454
- **`Tuple` element names** may differ if either the source or destination declares the tuple without named elements: an unnamed `Tuple` (e.g. `Tuple(Int32, Int32)`) is matched against the destination by element position and type only, not by name. For example, exporting from `t Tuple(Int32, Int32)` to `t Tuple(x Int32, y Int32)` is allowed as long as element types match positionally.
5555

56-
The following must match between source and destination:
56+
The following requirements apply to the source and destination:
5757

5858
1. **Column count** - source and destination must have the same number of columns by default. A mismatch in either direction throws `NUMBER_OF_COLUMNS_DOESNT_MATCH`. Set `export_merge_tree_part_schema_mismatch_mode = 'ignore_extra_source_columns_by_position'` to allow a source table with extra trailing columns; the destination having more columns than the source is still rejected in this mode.
59-
2. **`PARTITION BY` expressions** - for destinations other than data lakes, the source and destination `PARTITION BY` expressions must be identical. For Apache Iceberg destinations, the source partition key must be representable as an Iceberg partition spec and must match the destination partition fields and transforms.
59+
2. **`PARTITION BY` expressions** - the whole part must land in a single destination partition. Identical expressions always satisfy this; otherwise the destination expression has to be computable from the values the source partition key pins, or be proven single-valued over the part's min/max range. The same requirement applies to the partition fields and transforms of an Apache Iceberg destination. See [Source partition key compatibility](/docs/en/antalya/partition_export.md#source-partition-key-compatibility).
6060
3. **The position of every column backing the partition key** - it is not enough for the `PARTITION BY` expressions to be textually identical: every top-level column that provides a column or subcolumn used by the source table's partition key must have the same name at the same position in the destination table's schema. If such a column contains a named `Tuple`, its element names must also be declared in the same order (an unnamed `Tuple` on either side is exempt from this, per the allowance above). This comparison is recursive through nested tuples and through container types such as `Array` and `Map`.
6161

6262
For example, `CREATE TABLE src (a Int32, b Int32) ... PARTITION BY a` and `CREATE TABLE dst (b Int32, a Int32) ... PARTITION BY a` both have the expression `PARTITION BY a`, but `a` is at position 0 in `src` and position 1 in `dst`. The export is rejected with a `BAD_ARGUMENTS` exception whose message includes `Cannot export to <destination>: partition key column 'a' is at position 0 in the source table, but the destination's column at that position is named 'b'`.

‎docs/en/antalya/partition_export.md‎

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,13 @@ The manifest file produced by the commit contains a summary field `clickhouse.ex
2424

2525
The Iceberg manifest files contain statistics about the data. Exporting a merge tree partition is a non ephemeral long running task, in which nodes can be turned off and turned on. This means the stats of individual files need to be persisted somewhere in order to produce the final manifest. This is implemented through sidecars. Each data file exported will contain a "sibling" sidecar file named `<data_file_name>_clickhouse_export_part_sidecar.avro`. ClickHouse does not clean up these files, and they can be safely deleted once the data is comitted.
2626

27+
#### Source partition key compatibility
28+
29+
The source partition must not be split in the destination. This is validated at schedule time through two mechanisms:
30+
31+
1. Structural match: in case the source and destination are identical, the destination expression is a subset of the source expression or the destination expression can be entirely computed using only constants and the exact values guaranteed (pinned) by the source.
32+
2. Dynamic proof: the destination expression is monotonic over the source partition min/max range.
33+
2734
### On plain object storage exports:
2835

2936
Each MergeTree part will become a separate file with the following name convention: `<table_directory>/<partitioning>/<data_part_name>_<merge_tree_part_checksum>.<format>`. To ensure atomicity, a commit file containing the relative paths of all exported parts is also shipped. A data file should only be considered part of the dataset if a commit file references it. The commit file will be named using the following convention: `<table_directory>/commit_<partition_id>_<transaction_id>`.
@@ -45,10 +52,10 @@ TO TABLE [destination_database.]destination_table
4552

4653
## Requirements
4754

48-
`EXPORT PARTITION` exports each part via the same mechanism as [`EXPORT PART`](/docs/en/antalya/part_export.md#requirements), so the source and destination tables must satisfy the same compatibility requirements. Column names may differ (columns are matched by position, not by name), and column types may differ as long as they are safely castable (or `export_merge_tree_part_allow_lossy_cast = 1` is set). Beyond that, the following must match:
55+
`EXPORT PARTITION` exports each part via the same mechanism as [`EXPORT PART`](/docs/en/antalya/part_export.md#requirements), so the source and destination tables must satisfy the same compatibility requirements. Column names may differ (columns are matched by position, not by name), and column types may differ as long as they are safely castable (or `export_merge_tree_part_allow_lossy_cast = 1` is set). Beyond that, the following requirements apply:
4956

5057
1. **Column count** - source and destination must have the same number of columns by default. Set `export_merge_tree_part_schema_mismatch_mode = 'ignore_extra_source_columns_by_position'` to allow a source table with extra trailing columns; the destination having more columns than the source is still rejected in this mode.
51-
2. **`PARTITION BY` expressions** - for destinations other than data lakes, the source and destination `PARTITION BY` expressions must be identical. For Apache Iceberg destinations, the source partition key must match the destination partition fields and transforms.
58+
2. **`PARTITION BY` expressions** - the whole source partition must land in a single destination partition. Identical expressions always satisfy this; otherwise the destination expression has to be computable from the values the source partition key pins, or be proven single-valued over the partition's min/max range. The same requirement applies to the partition fields and transforms of an Apache Iceberg destination. See [Source partition key compatibility](#source-partition-key-compatibility).
5259
3. **Partition key column positions and layouts** - every top-level column that provides a column or subcolumn used by the source table's partition key must have the same name at the same position in the destination table's schema. Named `Tuple` elements within such a column must also be declared in the same order, including tuples nested inside `Array` or `Map`. This applies even if both tables' `PARTITION BY` expressions are textually identical. See [`EXPORT PART` requirements](/docs/en/antalya/part_export.md#requirements) for a worked example and the corresponding exception message.
5360

5461
## Settings

‎src/Core/FormatFactorySettings.h‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -199,7 +199,13 @@ When reading Parquet files, parse JSON columns as ClickHouse JSON Column.
199199
Schedule prefetches more aggressively if memory usage is below than threshold. Potentially useful e.g. if there are many small bloom filters to read over network.
200200
)", 0) \
201201
DECLARE(UInt64, input_format_parquet_memory_high_watermark, 4ul << 30, R"(
202-
Approximate memory limit for Parquet reader v3. Limits how many row groups or columns can be read in parallel. When reading multiple files in one query, the limit is on total memory usage across those files.
202+
Approximate memory limit for the Parquet reader. Limits how many row groups or columns can be read in parallel. When reading multiple files in one query, the limit is on total memory usage across those files.
203+
)", 0) \
204+
DECLARE(Double, input_format_parquet_prefetch_memory_fraction, 0.6, R"(
205+
Advanced tuning knob for the Parquet reader scheduler. Of the memory budget reserved for column data, the fraction given to compressed read-ahead (the `ColumnDataPrefetch` stage) versus decoded output (the `ColumnData` stage); the rest goes to decode. A higher value keeps more compressed pages in flight to hide read latency (useful on high-latency storage such as S3); a lower value caps read-ahead and leaves more budget for decoded columns. Must be in [0, 1]. The index and bloom-filter stages keep a fixed share of the memory budget regardless of this setting.
206+
)", 0) \
207+
DECLARE(Double, input_format_parquet_decode_thread_fraction, 0.375, R"(
208+
Advanced tuning knob for the Parquet reader scheduler. The fraction of the Parquet parsing thread pool dedicated to column decoding (the `ColumnData` stage); the remaining stages, which only issue asynchronous reads, share the rest. Raise it to give decoding (the only CPU-bound stage) more parallelism on fast/local storage; the default suits latency-bound remote reads where memory, not threads, limits concurrency. Must be in [0, 1].
203209
)", 0) \
204210
DECLARE(Bool, input_format_parquet_page_filter_push_down, true, R"(
205211
Skip pages using min/max values from column index.

‎src/Core/SettingsChangesHistory.cpp‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,8 @@ const VersionToSettingsChangesMap & getSettingsChangesHistory()
4949
{"analyzer_compatibility_apply_final_to_all_joined_tables", true, false, "Fixed a bug in the analyzer where FINAL on the left-most table of a JOIN was incorrectly applied to the other joined tables as well. previous_value=true so `compatibility` with versions before 26.6 restores the old behavior."},
5050
{"analyzer_compatibility_allow_non_aggregate_in_having", false, false, "New compatibility setting. When enabled, the new analyzer mimics the legacy `HAVING`-to-`WHERE` rewrite for non-aggregate AND-conjuncts instead of raising `NOT_AN_AGGREGATE`."},
5151
{"reserve_memory", 0, 0, "New setting to reserve memory for specific workload before starting a query."},
52+
{"input_format_parquet_prefetch_memory_fraction", 0.6, 0.6, "New setting to tune the Parquet reader split of the column-data memory budget between compressed read-ahead and decode."},
53+
{"input_format_parquet_decode_thread_fraction", 0.375, 0.375, "New setting to tune the Parquet reader share of the parsing thread pool given to column decoding."},
5254
{"output_format_image_width", 1024, 1024, "New setting controlling the width of the output image for image output formats such as PNG."},
5355
{"output_format_image_height", 1024, 1024, "New setting controlling the height of the output image for image output formats such as PNG."},
5456
{"output_format_image_terminal_mode", "", "", "New setting controlling whether image output formats such as PNG are rendered directly to the terminal using an inline image protocol."},

0 commit comments

Comments
 (0)