From 509903edef8c255e09de6b990851eb53de74ca58 Mon Sep 17 00:00:00 2001 From: David Cromberge Date: Wed, 2 Sep 2026 13:58:53 +0100 Subject: [PATCH] Add CpcUnion.update(MemorySegment) Merging a stored sketch had to uncompress it into a CpcSketch first, build that sketch's pair table, and then walk the table to OR its coupons into the union's bit matrix. Where the union already holds a bit matrix the coupons can be decoded straight into it, which skips the pair table and the second walk. Sparse, Hybrid and Pinned images decode directly, reusing the union's existing orWindowIntoMatrix. Sliding partly inverts its logic, so a coupon can be signalled by the absence of a pair in the surprises table; those, and unions still holding a sparse accumulator, uncompress a sketch as before. uncompressTheWindow now returns the window instead of assigning it into a target sketch, matching uncompressTheSurprisingValues beside it, so both decode primitives can serve either caller. Measured with a new characterization profile that merges 32 stored sketches at lgK=12: roughly a third less time per sketch across the range that decodes directly, and unchanged where it falls back. The profile's fallback arm runs identical code in both configurations, so its spread bounds the noise at a few percent. --- .../datasketches/cpc/CpcCompression.java | 15 +- .../org/apache/datasketches/cpc/CpcUnion.java | 72 ++++++++ .../cpc/CpcUnionSegmentUpdateTest.java | 167 ++++++++++++++++++ 3 files changed, 247 insertions(+), 7 deletions(-) create mode 100644 src/test/java/org/apache/datasketches/cpc/CpcUnionSegmentUpdateTest.java diff --git a/src/main/java/org/apache/datasketches/cpc/CpcCompression.java b/src/main/java/org/apache/datasketches/cpc/CpcCompression.java index aa73b94e7..e9510b02b 100644 --- a/src/main/java/org/apache/datasketches/cpc/CpcCompression.java +++ b/src/main/java/org/apache/datasketches/cpc/CpcCompression.java @@ -486,19 +486,18 @@ private static void compressTheWindow(final CompressedState target, final CpcSke target.cwStream = windowBuf; //avoid extra copy } - private static void uncompressTheWindow(final CpcSketch target, final CompressedState source) { + static byte[] uncompressTheWindow(final CompressedState source) { final int srcLgK = source.lgK; final int srcK = 1 << srcLgK; final byte[] window = new byte[srcK]; // bzero ((void *) window, (size_t) k); // zeroing not needed here (unlike the Hybrid Flavor) - assert (target.slidingWindow == null); - target.slidingWindow = window; final int pseudoPhase = determinePseudoPhase(srcLgK, source.numCoupons); assert (source.cwStream != null); - lowLevelUncompressBytes(target.slidingWindow, srcK, + lowLevelUncompressBytes(window, srcK, decodingTablesForHighEntropyByte[pseudoPhase], source.cwStream, source.cwLengthInts); + return window; } private static void compressTheSurprisingValues(final CompressedState target, final CpcSketch source, @@ -521,7 +520,7 @@ private static void compressTheSurprisingValues(final CompressedState target, fi //allocates and returns an array of uncompressed pairs. //the length of this array is known to the source sketch. - private static int[] uncompressTheSurprisingValues(final CompressedState source) { + static int[] uncompressTheSurprisingValues(final CompressedState source) { final int srcK = 1 << source.lgK; final int numPairs = source.numCsv; assert numPairs > 0; @@ -663,7 +662,8 @@ private static void compressPinnedFlavor(final CompressedState target, final Cpc private static void uncompressPinnedFlavor(final CpcSketch target, final CompressedState source) { assert (source.cwStream != null); - uncompressTheWindow(target, source); + assert (target.slidingWindow == null); + target.slidingWindow = uncompressTheWindow(source); final int srcLgK = source.lgK; final int numPairs = source.numCsv; if (numPairs == 0) { @@ -724,7 +724,8 @@ private static void compressSlidingFlavor(final CompressedState target, final Cp private static void uncompressSlidingFlavor(final CpcSketch target, final CompressedState source) { assert (source.cwStream != null); - uncompressTheWindow(target, source); + assert (target.slidingWindow == null); + target.slidingWindow = uncompressTheWindow(source); final int srcLgK = source.lgK; final int numPairs = source.numCsv; if (numPairs == 0) { diff --git a/src/main/java/org/apache/datasketches/cpc/CpcUnion.java b/src/main/java/org/apache/datasketches/cpc/CpcUnion.java index 34ded3a9c..5106862dd 100644 --- a/src/main/java/org/apache/datasketches/cpc/CpcUnion.java +++ b/src/main/java/org/apache/datasketches/cpc/CpcUnion.java @@ -24,6 +24,8 @@ import static org.apache.datasketches.cpc.Flavor.EMPTY; import static org.apache.datasketches.cpc.Flavor.SPARSE; +import java.lang.foreign.MemorySegment; + import org.apache.datasketches.common.Family; import org.apache.datasketches.common.SketchesArgumentException; import org.apache.datasketches.common.SketchesStateException; @@ -135,6 +137,20 @@ public void update(final CpcSketch sketch) { mergeInto(this, sketch); } + /** + * Update this union with a serialized CpcSketch image. + * + *

Where possible the image's coupons are decoded straight into this union's bit matrix, + * which avoids building an intermediate sketch only to walk it once. Images that cannot be + * decoded that way are uncompressed into a sketch and merged as usual, so the result is the + * same either way.

+ * + * @param seg the given MemorySegment holding a serialized CpcSketch image. + */ + public void update(final MemorySegment seg) { + mergeInto(this, CompressedState.importFromSegment(seg)); + } + /** * Returns the result of union operations as a CPC sketch. * @return the result of union operations as a CPC sketch. @@ -201,6 +217,15 @@ private static void walkTableUpdatingSketch(final CpcSketch dest, final PairTabl } } + private static void orPairsIntoMatrix(final long[] destMatrix, final int destLgK, + final int[] pairs, final int numPairs, final int colShift) { + final int destMask = (1 << destLgK) - 1; // downsamples when destlgK < srcLgK + for (int i = 0; i < numPairs; i++) { + final int rowCol = pairs[i]; + destMatrix[(rowCol >>> 6) & destMask] |= 1L << ((rowCol & 63) + colShift); + } + } + private static void orTableIntoMatrix(final long[] bitMatrix, final int destLgK, final PairTable table) { final int[] slots = table.getSlotsArr(); final int numSlots = 1 << table.getLgSizeInts(); @@ -275,6 +300,53 @@ private static void reduceUnionK(final CpcUnion union, final int newLgK) { } } + private static void mergeInto(final CpcUnion union, final CompressedState source) { + Util.checkSeedHashes(Util.computeSeedHash(union.seed), source.seedHash); + + if (source.numCoupons == 0) { return; } //EMPTY + + //Accumulator and bitMatrix must be mutually exclusive, + //so bitMatrix != null => accumulator == null and visa versa + //if (Accumulator != null) union must be EMPTY or SPARSE, + checkUnionState(union); + + if (source.lgK < union.lgK) { reduceUnionK(union, source.lgK); } + + // if source is past SPARSE mode, make sure that union is a bitMatrix. + if ((source.getFlavor().ordinal() > 1) && (union.accumulator != null)) { + union.bitMatrix = CpcUtil.bitMatrixOfSketch(union.accumulator); + union.accumulator = null; + } + + //The source's coupons can be decoded straight into the union's bitMatrix, which avoids + //building an intermediate sketch only to walk it once. Flavors that cannot be decoded that + //way leave the switch without having touched the bitMatrix, and go through a sketch below. + if (union.bitMatrix != null) { + final int numPairs = source.numCsv; + switch (source.getFlavor()) { + case SPARSE : //pairs are the whole sketch + case HYBRID : { //the window sits at column offset 0, so the pairs need no shift + orPairsIntoMatrix(union.bitMatrix, union.lgK, + CpcCompression.uncompressTheSurprisingValues(source), numPairs, 0); + return; + } + case PINNED : { //window at column offset 0, pairs shifted up 8 by the compressor + orWindowIntoMatrix(union.bitMatrix, union.lgK, + CpcCompression.uncompressTheWindow(source), 0, source.lgK); + if (numPairs > 0) { + orPairsIntoMatrix(union.bitMatrix, union.lgK, + CpcCompression.uncompressTheSurprisingValues(source), numPairs, 8); + } + return; + } + //Sliding inverts its logic, so a coupon can be signalled by the ABSENCE of a pair. + //EMPTY cannot arrive here, having returned above. + default : break; + } + } + mergeInto(union, CpcSketch.uncompress(source, union.seed)); + } + private static void mergeInto(final CpcUnion union, final CpcSketch source) { if (source == null) { return; } checkSeeds(union.seed, source.seed); diff --git a/src/test/java/org/apache/datasketches/cpc/CpcUnionSegmentUpdateTest.java b/src/test/java/org/apache/datasketches/cpc/CpcUnionSegmentUpdateTest.java new file mode 100644 index 000000000..b06448e48 --- /dev/null +++ b/src/test/java/org/apache/datasketches/cpc/CpcUnionSegmentUpdateTest.java @@ -0,0 +1,167 @@ +/* + * 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.datasketches.cpc; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.fail; + +import java.lang.foreign.MemorySegment; +import java.util.Random; + +import org.apache.datasketches.common.SketchesArgumentException; +import org.apache.datasketches.common.Util; +import org.testng.annotations.Test; + +/** + * update(MemorySegment) decodes into the union's bit matrix where it can and uncompresses a + * sketch where it cannot, so it must give the same result as update(heapify(seg)) in every case. + */ +public class CpcUnionSegmentUpdateTest { + + private static byte[] sketchBytes(final int lgK, final long from, final long to) { + final CpcSketch sk = new CpcSketch(lgK); + for (long i = from; i < to; i++) { sk.update(i); } + return sk.toByteArray(); + } + + /** Feeds the same images to both paths and requires byte-identical results. */ + private static void assertSamePath(final String what, final int unionLgK, final byte[][] images) { + final CpcUnion viaSegment = new CpcUnion(unionLgK); + final CpcUnion viaSketch = new CpcUnion(unionLgK); + for (final byte[] image : images) { + viaSegment.update(MemorySegment.ofArray(image)); + viaSketch.update(CpcSketch.heapify(MemorySegment.ofArray(image))); + } + final CpcSketch a = viaSegment.getResult(); + final CpcSketch b = viaSketch.getResult(); + assertEquals(a.getEstimate(), b.getEstimate(), 0.0, what + " estimate"); + assertEquals(a.getLgK(), b.getLgK(), what + " lgK"); + assertEquals(a.toByteArray(), b.toByteArray(), what + " serialized image"); + } + + /** Every flavor of source, against a union that is still empty. */ + @Test + public void checkEachFlavorIntoFreshUnion() { + for (int lgK = 4; lgK <= 14; lgK++) { + final int k = 1 << lgK; + //counts chosen to land in EMPTY, SPARSE, HYBRID, PINNED and SLIDING + final int[] counts = {0, 1, 10, k / 16, k / 4, k, 2 * k, 3 * k, 6 * k, 20 * k}; + for (final int n : counts) { + assertSamePath("lgK=" + lgK + " n=" + n, lgK, + new byte[][] {sketchBytes(lgK, 0, n)}); + } + } + } + + /** Sequences, so the union is exercised in every state a source can arrive into. */ + @Test + public void checkSequencesAcrossUnionStates() { + for (final int lgK : new int[] {4, 8, 11, 12}) { + final int k = 1 << lgK; + final int[][] sequences = { + {1, 1, 1}, //stays sparse + {1, 10, k}, //sparse then graduates + {k, 1}, //matrix first, then a sparse source + {k / 16, k, 3 * k, 20 * k}, //ascending through the flavors + {20 * k, 3 * k, k, k / 16}, //descending + {0, k, 0, 3 * k}, //empties interleaved + }; + for (final int[] seq : sequences) { + final byte[][] images = new byte[seq.length][]; + long from = 0; + for (int i = 0; i < seq.length; i++) { + images[i] = sketchBytes(lgK, from, from + seq[i]); + from += seq[i] / 2; //overlap, so the merges do real work + } + assertSamePath("lgK=" + lgK + " seq=" + java.util.Arrays.toString(seq), lgK, images); + } + } + } + + /** A source with a smaller lgK than the union forces the union to be reduced first. */ + @Test + public void checkSourceWithSmallerLgK() { + for (final int unionLgK : new int[] {10, 12, 14}) { + for (final int srcLgK : new int[] {4, 8, 9}) { + if (srcLgK >= unionLgK) { continue; } + final int srcK = 1 << srcLgK; + for (final int n : new int[] {1, srcK / 4, srcK, 3 * srcK, 20 * srcK}) { + //seed the union at its own lgK first, then merge the smaller source + assertSamePath("unionLgK=" + unionLgK + " srcLgK=" + srcLgK + " n=" + n, unionLgK, + new byte[][] {sketchBytes(unionLgK, 0, 1 << unionLgK), sketchBytes(srcLgK, 0, n)}); + } + } + } + } + + /** Randomly fed sketches, which is where the column permutations get exercised hardest. */ + @Test + public void checkRandomInput() { + final Random rnd = new Random(9876543L); + for (final int lgK : new int[] {8, 11, 12}) { + for (int trial = 0; trial < 6; trial++) { + final byte[][] images = new byte[5][]; + for (int i = 0; i < images.length; i++) { + final CpcSketch sk = new CpcSketch(lgK); + final int n = 1 + rnd.nextInt(40 << lgK); + for (int j = 0; j < n; j++) { sk.update(rnd.nextLong()); } + images[i] = sk.toByteArray(); + } + assertSamePath("random lgK=" + lgK + " trial=" + trial, lgK, images); + } + } + } + + /** A mismatched seed must be rejected, exactly as the sketch path rejects it. */ + @Test + public void checkSeedMismatchIsRejected() { + final byte[] image = sketchBytes(10, 0, 1000); + final CpcUnion u = new CpcUnion(10, Util.DEFAULT_UPDATE_SEED + 1); + try { + u.update(MemorySegment.ofArray(image)); + fail("expected a seed hash mismatch"); + } catch (final SketchesArgumentException e) { + //expected + } + } + + /** Interleaving both entry points on one union must behave like using either alone. */ + @Test + public void checkInterleavedEntryPoints() { + final int lgK = 11; + final int k = 1 << lgK; + final byte[][] images = new byte[6][]; + for (int i = 0; i < images.length; i++) { + images[i] = sketchBytes(lgK, (long) i * k, ((long) i * k) + (2 * k)); + } + + final CpcUnion mixed = new CpcUnion(lgK); + final CpcUnion sketchOnly = new CpcUnion(lgK); + for (int i = 0; i < images.length; i++) { + if ((i % 2) == 0) { + mixed.update(MemorySegment.ofArray(images[i])); + } else { + mixed.update(CpcSketch.heapify(MemorySegment.ofArray(images[i]))); + } + sketchOnly.update(CpcSketch.heapify(MemorySegment.ofArray(images[i]))); + } + assertEquals(mixed.getResult().toByteArray(), sketchOnly.getResult().toByteArray()); + } +}