voonhous commented on code in PR #19707:
URL: https://github.com/apache/hudi/pull/19707#discussion_r3833034380
##########
hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestShardRecordIterator.java:
##########
@@ -361,4 +367,70 @@ void testFirstPageThrottleRetriesAndSucceeds() {
assertEquals(200, captor.getAllValues().get(0).limit());
assertEquals(100, captor.getAllValues().get(1).limit());
}
+
+ // -------------------------------------------------------------------------
+ // KPL de-aggregation
+ // -------------------------------------------------------------------------
Review Comment:
Extracted the constant and the frame builders into `KplTestUtils`; both test
classes now share them.
##########
hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/helpers/KinesisDeaggregator.java:
##########
@@ -44,36 +59,137 @@ public static List<Record> deaggregate(List<Record>
records) {
if (records == null || records.isEmpty()) {
return new ArrayList<>();
}
- List<com.amazonaws.services.kinesis.model.Record> v1Records = new
ArrayList<>(records.size());
- for (Record r : records) {
- v1Records.add(toV1Record(r));
- }
- List<UserRecord> userRecords = UserRecord.deaggregate(v1Records);
- List<Record> result = new ArrayList<>(userRecords.size());
- for (UserRecord ur : userRecords) {
- result.add(toV2Record(ur));
+ List<Record> result = new ArrayList<>(records.size());
+ for (Record record : records) {
+ byte[] data = record.data() == null ? null : record.data().asByteArray();
+ if (!isAggregated(data)) {
+ result.add(record);
+ continue;
+ }
+ int payloadLength = data.length - MAGIC.length - DIGEST_LENGTH;
+ try {
+ result.addAll(expand(record, data, MAGIC.length, payloadLength));
+ } catch (IOException e) {
+ log.warn("Kinesis record with sequence number {} has a matching KPL
aggregation digest but could not be "
+ + "decoded, passing it through as-is", record.sequenceNumber(), e);
+ result.add(record);
+ }
}
return result;
}
- private static com.amazonaws.services.kinesis.model.Record toV1Record(Record
v2) {
- com.amazonaws.services.kinesis.model.Record v1 = new
com.amazonaws.services.kinesis.model.Record();
- v1.withData(ByteBuffer.wrap(v2.data().asByteArray()));
- v1.withPartitionKey(v2.partitionKey());
- v1.withSequenceNumber(v2.sequenceNumber());
- if (v2.approximateArrivalTimestamp() != null) {
-
v1.withApproximateArrivalTimestamp(Date.from(v2.approximateArrivalTimestamp()));
+ private static boolean isAggregated(byte[] data) {
+ // Strictly greater: a frame with an empty payload is not an aggregate,
matching KCL.
+ if (data == null || data.length <= MAGIC.length + DIGEST_LENGTH) {
+ return false;
+ }
+ for (int i = 0; i < MAGIC.length; i++) {
+ if (data[i] != MAGIC[i]) {
+ return false;
+ }
}
- return v1;
+ byte[] expectedDigest = new byte[DIGEST_LENGTH];
+ System.arraycopy(data, data.length - DIGEST_LENGTH, expectedDigest, 0,
DIGEST_LENGTH);
+ byte[] actualDigest = md5(data, MAGIC.length, data.length - MAGIC.length -
DIGEST_LENGTH);
+ return MessageDigest.isEqual(actualDigest, expectedDigest);
}
- private static Record toV2Record(UserRecord v1) {
+ private static byte[] md5(byte[] data, int offset, int length) {
+ try {
+ MessageDigest digest = MessageDigest.getInstance("MD5");
+ digest.update(data, offset, length);
+ return digest.digest();
+ } catch (NoSuchAlgorithmException e) {
+ throw new IllegalStateException("MD5 is not available in this JVM", e);
+ }
+ }
+
+ /**
+ * Parses the {@code AggregatedRecord} message: repeated string
partition_key_table (field 1),
+ * repeated string explicit_hash_key_table (field 2) and repeated Record
records (field 3).
+ */
+ private static List<Record> expand(Record parent, byte[] data, int offset,
int length) throws IOException {
+ CodedInputStream input = CodedInputStream.newInstance(data, offset,
length);
+ List<String> partitionKeyTable = new ArrayList<>();
+ List<String> explicitHashKeyTable = new ArrayList<>();
+ List<byte[]> subMessages = new ArrayList<>();
+ while (!input.isAtEnd()) {
+ int tag = input.readTag();
+ switch (tag) {
+ case 10:
+ partitionKeyTable.add(input.readStringRequireUtf8());
Review Comment:
Added the tag formula to both parser javadocs and a `// field N, wire type`
comment on each case.
--
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]