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]

Reply via email to