Copilot commented on code in PR #2623: URL: https://github.com/apache/phoenix/pull/2623#discussion_r4099357997
########## phoenix-core-server/src/main/java/org/apache/phoenix/mapreduce/RegionServerSplitCoalescer.java: ########## @@ -0,0 +1,183 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.phoenix.mapreduce; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.apache.hadoop.mapreduce.InputSplit; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Coalesces region-boundary {@link PhoenixInputSplit}s into one split per RegionServer, reducing + * mapper fan-out and hot spotting on large (e.g. salted) tables. Grouping uses the + * {@code host:port} identity already stamped on each split + * ({@link PhoenixInputSplit#getRegionServerName()}) — no region-location RPCs, and the port keeps + * two RegionServers on one host separate. A split with no identity falls back to its hostname + * ({@link InputSplit#getLocations()}), else {@link #UNKNOWN_SERVER}. + */ +final class RegionServerSplitCoalescer { + + private static final Logger LOGGER = LoggerFactory.getLogger(RegionServerSplitCoalescer.class); + + /** + * Sentinel key for splits with no location (empty {@link InputSplit#getLocations()}, e.g. a + * region-in-transition). They are coalesced together rather than failing the job, since + * coalescing is an optimisation, not a correctness requirement. + */ + static final String UNKNOWN_SERVER = "UNKNOWN_SERVER"; + + private RegionServerSplitCoalescer() { + } + + /** + * Coalesces the given splits by RegionServer, guarding correctness. Since coalescing must never + * change which rows are processed, this returns the original splits unchanged when there is + * nothing to coalesce ({@code <= 1} split or {@code null}), when the scan count is not preserved, + * or when coalescing throws. {@link InterruptedException} is propagated (interrupt flag + * restored). + * @param splits region-granular splits to coalesce; may be {@code null} + * @return the coalesced splits, or {@code splits} unchanged when coalescing is skipped or + * rejected + */ + static List<InputSplit> coalesceWithGuard(List<InputSplit> splits) throws InterruptedException { + if (splits == null || splits.size() <= 1) { + return splits; + } + try { + List<InputSplit> coalesced = coalesce(splits); + if (!scanCountPreserved(splits, coalesced)) { + LOGGER.error( + "Split coalescing changed the scan count ({} -> {}); falling back to base splits to " + + "preserve correctness", + countScans(splits), countScans(coalesced)); + return splits; + } + LOGGER.info("Split coalescing: {} base splits coalesced into {} splits", splits.size(), + coalesced.size()); + return coalesced; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw e; + } catch (Exception e) { + LOGGER.error("Split coalescing failed; falling back to base splits", e); + return splits; + } + } + + /** + * Coalesces region-boundary splits so all regions on the same RegionServer become one + * {@link PhoenixInputSplit}, with each group's splits sorted by start key and their scans + * concatenated. Performs no correctness guard, unlike {@link #coalesceWithGuard(List)}; prefer + * the guarded entry point. + * @param splits region-granular splits to coalesce + * @return one coalesced split per distinct server location + */ + static List<InputSplit> coalesce(List<InputSplit> splits) + throws IOException, InterruptedException { + Map<String, List<PhoenixInputSplit>> splitsByServer = groupSplitsByServer(splits); + List<InputSplit> coalescedSplits = new ArrayList<>(splitsByServer.size()); + for (Map.Entry<String, List<PhoenixInputSplit>> entry : splitsByServer.entrySet()) { + List<PhoenixInputSplit> serverSplits = entry.getValue(); + // Sort by start key so each mapper scans its server's regions in key order. + serverSplits.sort((s1, s2) -> Bytes.compareTo(s1.getKeyRange().getLowerRange(), + s2.getKeyRange().getLowerRange())); + coalescedSplits.add(createCoalescedSplit(serverSplits, entry.getKey())); + } + return coalescedSplits; + } + + /** + * Groups splits by RegionServer identity: the {@code host:port} stamped on each split + * ({@link PhoenixInputSplit#getRegionServerName()}), else the bare hostname + * ({@link InputSplit#getLocations()}), else {@link #UNKNOWN_SERVER}. Uses a {@link LinkedHashMap} + * so grouping order is deterministic (first-seen first) given a fixed input order. + */ + private static Map<String, List<PhoenixInputSplit>> groupSplitsByServer(List<InputSplit> splits) + throws IOException, InterruptedException { + Map<String, List<PhoenixInputSplit>> splitsByServer = new LinkedHashMap<>(); + for (InputSplit split : splits) { + PhoenixInputSplit pSplit = (PhoenixInputSplit) split; + String serverName = pSplit.getRegionServerName(); + String serverKey; + if (serverName != null && !serverName.isEmpty()) { + serverKey = serverName; + } else { + String[] locations = pSplit.getLocations(); + if (locations != null && locations.length > 0 && locations[0] != null) { + serverKey = locations[0]; + LOGGER.warn("Split {} has no RegionServer identity; grouping by hostname {} instead", + Bytes.toStringBinary(pSplit.getKeyRange().getLowerRange()), serverKey); + } else { + serverKey = UNKNOWN_SERVER; + LOGGER.warn("Split {} has no location (region may be in transition); assigning to {}", + Bytes.toStringBinary(pSplit.getKeyRange().getLowerRange()), UNKNOWN_SERVER); + } + } + splitsByServer.computeIfAbsent(serverKey, k -> new ArrayList<>()).add(pSplit); + } + return splitsByServer; + } + + /** + * Creates one coalesced {@link PhoenixInputSplit} by concatenating the scans of the given + * per-region splits (already sorted by start key) and summing their sizes. It carries + * {@code serverKey} as its RegionServer identity and the group's shared hostname (taken from the + * first member) as its data-locality location. + */ + private static PhoenixInputSplit createCoalescedSplit(List<PhoenixInputSplit> splits, + String serverKey) throws IOException, InterruptedException { + List<Scan> allScans = new ArrayList<>(); + long totalSize = 0; + for (PhoenixInputSplit split : splits) { + allScans.addAll(split.getScans()); + totalSize += split.getLength(); + } + String[] firstLocations = splits.get(0).getLocations(); + String hostname = firstLocations.length > 0 ? firstLocations[0] : null; + if (LOGGER.isDebugEnabled()) { + LOGGER.debug("Created coalesced split with {} regions from server {}", splits.size(), + serverKey); + } + return new PhoenixInputSplit(allScans, totalSize, hostname, serverKey); + } + + /** + * Whether coalescing preserved every region scan (none dropped or duplicated). The guard + * {@link #coalesceWithGuard(List)} uses to decide whether the coalesced result is safe. + */ + static boolean scanCountPreserved(List<InputSplit> base, List<InputSplit> coalesced) { + return countScans(base) == countScans(coalesced); Review Comment: This guard only compares totals, so a result that drops one scan and duplicates another still passes even though it changes the rows processed. Since this method is the advertised correctness fallback, compare a multiset of the original scans (using identity or a complete scan fingerprint) rather than only the count before accepting the coalesced result. ########## phoenix-core/src/it/java/org/apache/phoenix/mapreduce/PhoenixInputFormatSplitCoalescingIT.java: ########## @@ -0,0 +1,207 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.phoenix.mapreduce; + +import static org.apache.phoenix.util.TestUtil.TEST_PROPERTIES; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; + +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.util.HashSet; +import java.util.List; +import java.util.Properties; +import java.util.Set; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.util.Bytes; +import org.apache.hadoop.mapreduce.InputSplit; +import org.apache.hadoop.mapreduce.Job; +import org.apache.hadoop.mapreduce.lib.db.DBWritable; +import org.apache.phoenix.end2end.NeedsOwnMiniClusterTest; +import org.apache.phoenix.mapreduce.util.PhoenixConfigurationUtil; +import org.apache.phoenix.mapreduce.util.PhoenixMapReduceUtil; +import org.apache.phoenix.query.BaseTest; +import org.apache.phoenix.util.PropertiesUtil; +import org.apache.phoenix.util.ReadOnlyProps; +import org.junit.BeforeClass; +import org.junit.Test; +import org.junit.experimental.categories.Category; + +/** + * Integration test for RegionServer-level split coalescing in {@link PhoenixInputFormat} on a + * multi-RegionServer mini-cluster. Exercises what the {@code RegionServerSplitCoalescerTest} unit + * tests cannot: real region locations resolved from a live cluster. Verifies that enabling + * {@link PhoenixInputFormat#SPLIT_COALESCING_ENABLED} collapses the region-granular splits into one + * split per distinct RegionServer ({@code host:port}) while preserving every underlying region scan + * (no rows dropped or duplicated), and that the default (disabled) path leaves the region-granular + * splits untouched. + * <p> + * Note: a mini-cluster's RegionServers all share the same hostname and differ only by port, so + * grouping by {@code host:port} ({@link PhoenixInputSplit#getRegionServerName()}) — rather than by + * hostname — is exactly what keeps them in separate splits here. The test asserts against the + * dynamically-computed count of distinct RegionServer identities rather than a hard-coded count. + */ +@Category(NeedsOwnMiniClusterTest.class) +public class PhoenixInputFormatSplitCoalescingIT extends BaseTest { + + private static final int SALT_BUCKETS = 4; + private static String tableName; + + @BeforeClass + public static synchronized void doSetup() throws Exception { + NUM_SLAVES_BASE = 2; + setUpTestDriver(ReadOnlyProps.EMPTY_PROPS, ReadOnlyProps.EMPTY_PROPS); + createAndPopulateTable(); + } + + /** + * Salting pre-splits the table into {@code SALT_BUCKETS} regions at creation, giving more than + * one region-granular split (spread across the RegionServers) for coalescing to collapse. + */ + private static void createAndPopulateTable() throws Exception { + tableName = generateUniqueName(); + Properties props = PropertiesUtil.deepCopy(TEST_PROPERTIES); + try (Connection conn = DriverManager.getConnection(getUrl(), props)) { + conn.createStatement().execute("CREATE TABLE " + tableName + + " (PK INTEGER NOT NULL PRIMARY KEY, V VARCHAR) SALT_BUCKETS=" + SALT_BUCKETS); + try (PreparedStatement upsert = + conn.prepareStatement("UPSERT INTO " + tableName + " (PK, V) VALUES (?, ?)")) { + for (int i = 0; i < 100; i++) { + upsert.setInt(1, i); + upsert.setString(2, "v" + i); + upsert.executeUpdate(); + } + } + conn.commit(); + } + } + + @Test + public void coalescingCollapsesRegionSplitsToOnePerRegionServer() throws Exception { + List<InputSplit> baseline = getSplits(false); + List<InputSplit> coalesced = getSplits(true); + + assertTrue( + "Salted table should generate more than one region-granular split, got " + baseline.size(), + baseline.size() > 1); + + // Every region-granular split carries the host:port identity of the RegionServer hosting it. + Set<String> servers = new HashSet<>(); + for (InputSplit split : baseline) { + String server = ((PhoenixInputSplit) split).getRegionServerName(); + assertNotNull("Region split should have a RegionServer identity", server); + servers.add(server); + } + + assertEquals("Coalescing must produce exactly one split per distinct RegionServer (host:port)", + servers.size(), coalesced.size()); + assertTrue("Coalescing must reduce the split count (coalesced " + coalesced.size() + + " < baseline " + baseline.size() + ")", coalesced.size() < baseline.size()); + + // Directly assert coalescing happened, not just that the count dropped: with stats-splitting + // forced off each region is one scan, so a coalesced split carrying more than one scan means + // multiple regions were actually merged into a single mapper. + boolean merged = + coalesced.stream().anyMatch(s -> ((PhoenixInputSplit) s).getScans().size() > 1); + assertTrue("At least one coalesced split must merge multiple regions' scans", merged); + + for (InputSplit split : coalesced) { + assertNotEquals("Coalesced split must resolve to a real RegionServer", + RegionServerSplitCoalescer.UNKNOWN_SERVER, + ((PhoenixInputSplit) split).getRegionServerName()); + } + + // The count/location asserts above would still pass if a whole region's scan went missing, so + // assert coalescing preserved every region scan (this is what protects a downstream delete + // job). + assertEquals("Coalescing must preserve the total underlying scan count (no rows dropped)", + countScans(baseline), countScans(coalesced)); + assertEquals("Coalescing must preserve the exact set of region scan ranges", + scanRanges(baseline), scanRanges(coalesced)); + } + + @Test + public void coalescingDisabledYieldsOneSplitPerRegion() throws Exception { + // The feature is opt-in: with the flag off, getSplits must return the raw region-granular + // splits, i.e. exactly one split per HBase region. Asserting against the live region count is a + // direct expression of "no coalescing happened" that does not depend on how many scans a region + // maps to (unlike PhoenixInputSplit#isCoalesced, which is really "this split has > 1 scan"). + int regionCount = getUtility().getAdmin().getRegions(TableName.valueOf(tableName)).size(); + assertTrue("Salted table should have more than one region", regionCount > 1); + + List<InputSplit> baseline = getSplits(false); + assertEquals("With coalescing disabled, getSplits must produce one split per region", + regionCount, baseline.size()); + } + + /** + * Runs {@link PhoenixInputFormat#getSplits} against the live cluster with coalescing on/off. A + * {@link Job} is a {@link org.apache.hadoop.mapreduce.JobContext}, which is what + * {@code getSplits} takes. Stats-based splitting is forced off on both paths so the baseline and + * coalesced runs share the same region-granular starting point; otherwise the coalescing-off + * baseline would default to stats splitting ({@code DEFAULT_SPLIT_BY_STATS} is true) and the + * scan-set comparison would measure that difference rather than coalescing. + */ + private List<InputSplit> getSplits(boolean coalescingEnabled) throws Exception { + Configuration conf = new Configuration(getUtility().getConfiguration()); + Job job = Job.getInstance(conf); + PhoenixMapReduceUtil.setInput(job, DummyDBWritable.class, tableName, null, "PK"); + PhoenixConfigurationUtil.setSplitByStats(job.getConfiguration(), false); Review Comment: This pre-disables stats splitting for the enabled path, so the integration test never exercises the new behavior in `PhoenixInputFormat.getSplits()` that must turn stats splitting off when coalescing is requested. Start the enabled path with stats splitting on while keeping the baseline off; then the existing split and scan assertions will cover that requirement. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
