diff --git a/gluten-flink/ut/pom.xml b/gluten-flink/ut/pom.xml index c97c3838e12..3e3b9b5d8c8 100644 --- a/gluten-flink/ut/pom.xml +++ b/gluten-flink/ut/pom.xml @@ -281,6 +281,24 @@ ${flink.version} test + + org.apache.flink + flink-orc + ${flink.version} + test + + + com.google.protobuf + protobuf-java + 3.25.5 + test + + + org.apache.hadoop + hadoop-client + 2.7.2 + test + diff --git a/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java b/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java index 14b250480bc..828e5c56bc7 100644 --- a/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java +++ b/gluten-flink/ut/src/test/java/org/apache/gluten/table/runtime/stream/custom/NexmarkTest.java @@ -27,6 +27,12 @@ import com.salesforce.kafka.test.junit5.SharedKafkaTestResource; import com.salesforce.kafka.test.listeners.PlainListener; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch; +import org.apache.orc.OrcFile; +import org.apache.orc.Reader; +import org.apache.orc.RecordReader; +import org.apache.orc.TypeDescription; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; @@ -42,6 +48,7 @@ import java.nio.file.Paths; import java.util.ArrayList; import java.util.Arrays; +import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -147,32 +154,40 @@ void testKafkaSourceSqlPushesDownWatermark() { @Test void testAllNexmarkSourceQueries() throws ExecutionException, InterruptedException, TimeoutException { - setupNexmarkEnvironment(tEnv, "ddl_gen.sql", NEXMARK_VARIABLES); - List queryFiles = getQueries(); - assertThat(queryFiles).isNotEmpty(); - LOG.warn("Found {} Nexmark query files: {}", queryFiles.size(), queryFiles); - - for (String queryFile : queryFiles) { - LOG.warn("Executing nextmark query from file: {}", queryFile); - executeQuery(tEnv, queryFile, false); + try { + setupNexmarkEnvironment(tEnv, "ddl_gen.sql", NEXMARK_VARIABLES); + List queryFiles = getQueries(); + assertThat(queryFiles).isNotEmpty(); + LOG.warn("Found {} Nexmark query files: {}", queryFiles.size(), queryFiles); + + for (String queryFile : queryFiles) { + LOG.warn("Executing nextmark query from file: {}", queryFile); + executeQuery(tEnv, queryFile, false); + } + } finally { + clearEnvironment(tEnv); } - clearEnvironment(tEnv); } @Test void testAllKafkaSourceQueries() throws ExecutionException, InterruptedException, TimeoutException { - kafkaInstance.getKafkaTestUtils().createTopic(topicName, 1, (short) 1); - setupNexmarkEnvironment(tEnv, "ddl_kafka.sql", KAFKA_VARIABLES); - List queryFiles = getQueries(); - assertThat(queryFiles).isNotEmpty(); - LOG.warn("Found {} Nexmark query files: {}", queryFiles.size(), queryFiles); - - for (String queryFile : queryFiles) { - LOG.warn("Executing kafka query from file:{}", queryFile); - executeQuery(tEnv, queryFile, true); + try { + kafkaInstance.getKafkaTestUtils().createTopic(topicName, 1, (short) 1); + setupNexmarkEnvironment(tEnv, "ddl_kafka.sql", KAFKA_VARIABLES); + List queryFiles = getQueries(); + assertThat(queryFiles).isNotEmpty(); + LOG.warn("Found {} Nexmark query files: {}", queryFiles.size(), queryFiles); + + for (String queryFile : queryFiles) { + LOG.warn("Executing kafka query from file:{}", queryFile); + if (!"q10_orc.sql".equals(queryFile)) { + executeQuery(tEnv, queryFile, true); + } + } + } finally { + clearEnvironment(tEnv); } - clearEnvironment(tEnv); } private static void setupNexmarkEnvironment( @@ -206,7 +221,10 @@ private static void clearEnvironment(StreamTableEnvironment tEnv) { String sql = String.format("drop table if exists %s", tableName); tEnv.executeSql(sql); } + tEnv.executeSql("drop table if exists nexmark_q10_orc"); for (String view : VIEWS) { + String dropTemporaryViewSql = String.format("drop temporary view if exists %s", view); + tEnv.executeSql(dropTemporaryViewSql); String sql = String.format("drop view if exists %s", view); tEnv.executeSql(sql); } @@ -220,7 +238,15 @@ private static void clearEnvironment(StreamTableEnvironment tEnv) { private void executeQuery(StreamTableEnvironment tEnv, String queryFileName, boolean kafkaSource) throws ExecutionException, InterruptedException, TimeoutException { + if ("q10_orc.sql".equals(queryFileName) && !kafkaSource) { + executeQ10OrcBatchQuery(); + return; + } + String queryContent = readSqlFromFile(NEXMARK_RESOURCE_DIR + "/" + queryFileName); + if ("q10_orc.sql".equals(queryFileName)) { + cleanQ10OrcOutput(); + } String[] sqlStatements = queryContent.split(";"); assertThat(sqlStatements.length).isGreaterThanOrEqualTo(2); @@ -242,11 +268,142 @@ private void executeQuery(StreamTableEnvironment tEnv, String queryFileName, boo assertThat(checkJobRunningStatus(insertResult, 30000) == true); } else { waitForJobCompletion(insertResult, 30000); + if ("q10_orc.sql".equals(queryFileName)) { + verifyQ10OrcOutput(); + } } } assertTrue(sqlStatements[sqlStatements.length - 1].trim().isEmpty()); } + private void executeQ10OrcBatchQuery() + throws ExecutionException, InterruptedException, TimeoutException { + cleanQ10OrcOutput(); + StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); + env.setParallelism(1); + + EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build(); + StreamTableEnvironment batchTEnv = StreamTableEnvironment.create(env, settings); + try { + createQ10OrcBidView(batchTEnv); + String queryContent = readSqlFromFile(NEXMARK_RESOURCE_DIR + "/q10_orc.sql"); + String[] sqlStatements = queryContent.split(";"); + assertThat(sqlStatements.length).isEqualTo(3); + + TableResult createResult = batchTEnv.executeSql(sqlStatements[0].trim()); + assertFalse(createResult.getJobClient().isPresent()); + + TableResult insertResult = batchTEnv.executeSql(sqlStatements[1].trim()); + waitForJobCompletion(insertResult, 30000); + verifyQ10OrcOutput(); + assertTrue(sqlStatements[2].trim().isEmpty()); + } finally { + clearEnvironment(batchTEnv); + } + } + + private static void createQ10OrcBidView(StreamTableEnvironment tEnv) { + tEnv.executeSql( + "CREATE TEMPORARY VIEW bid AS " + + "SELECT " + + "CAST(1 AS BIGINT) AS auction, " + + "CAST(2 AS BIGINT) AS bidder, " + + "CAST(100 AS BIGINT) AS price, " + + "CAST('channel' AS STRING) AS channel, " + + "CAST('url' AS STRING) AS url, " + + "TIMESTAMP '2026-07-21 07:30:00' AS `dateTime`, " + + "CAST('extra' AS STRING) AS extra"); + } + + private void cleanQ10OrcOutput() { + Path outputDir = Paths.get("/tmp/data/output/bid_orc"); + if (!Files.exists(outputDir)) { + return; + } + try (java.util.stream.Stream files = Files.walk(outputDir)) { + files + .sorted(Comparator.reverseOrder()) + .forEach( + path -> { + try { + Files.deleteIfExists(path); + } catch (IOException e) { + throw new RuntimeException("Failed to delete " + path, e); + } + }); + } catch (IOException e) { + throw new RuntimeException("Failed to clean Q10 ORC output directory", e); + } + } + + private void verifyQ10OrcOutput() throws InterruptedException { + Path outputDir = Paths.get("/tmp/data/output/bid_orc"); + assertTrue("Q10 ORC output directory should exist", Files.exists(outputDir)); + + List partFiles = waitForFinalQ10OrcPartFiles(outputDir); + long rowCount = 0L; + for (Path partFile : partFiles) { + try { + rowCount += readAndVerifyQ10OrcFile(partFile); + } catch (IOException e) { + throw new RuntimeException("Failed to read Q10 ORC output file " + partFile, e); + } + } + assertThat(rowCount).isGreaterThan(0L); + } + + private List waitForFinalQ10OrcPartFiles(Path outputDir) throws InterruptedException { + long deadlineMillis = System.currentTimeMillis() + 30000L; + List regularFiles = List.of(); + while (System.currentTimeMillis() < deadlineMillis) { + try (java.util.stream.Stream files = Files.walk(outputDir)) { + regularFiles = files.filter(Files::isRegularFile).sorted().collect(Collectors.toList()); + } catch (IOException e) { + throw new RuntimeException("Failed to inspect Q10 ORC output", e); + } + + boolean hasInProgress = + regularFiles.stream().anyMatch(path -> path.toString().contains(".inprogress")); + List partFiles = + regularFiles.stream() + .filter(path -> path.getFileName().toString().startsWith("part-")) + .collect(Collectors.toList()); + if (!hasInProgress && !partFiles.isEmpty()) { + return partFiles; + } + Thread.sleep(1000L); + } + + assertThat(regularFiles).allMatch(path -> !path.toString().contains(".inprogress")); + List partFiles = + regularFiles.stream() + .filter(path -> path.getFileName().toString().startsWith("part-")) + .collect(Collectors.toList()); + assertThat(partFiles).isNotEmpty(); + return partFiles; + } + + private long readAndVerifyQ10OrcFile(Path partFile) throws IOException { + Reader reader = + OrcFile.createReader( + new org.apache.hadoop.fs.Path(partFile.toUri()), + OrcFile.readerOptions(new Configuration())); + TypeDescription schema = reader.getSchema(); + assertThat(schema.getCategory()).isEqualTo(TypeDescription.Category.STRUCT); + assertThat(schema.getFieldNames()) + .containsExactly("auction", "bidder", "price", "dateTime", "extra"); + + long rowCount = 0L; + try (RecordReader rows = reader.rows()) { + VectorizedRowBatch batch = schema.createRowBatch(); + while (rows.nextBatch(batch)) { + rowCount += batch.size; + } + } + assertThat(rowCount).isEqualTo(reader.getNumberOfRows()); + return rowCount; + } + private void waitForJobCompletion(TableResult result, long timeoutMs) throws InterruptedException, ExecutionException, TimeoutException { assertTrue(result.getJobClient().isPresent()); diff --git a/gluten-flink/ut/src/test/resources/nexmark/q10_orc.sql b/gluten-flink/ut/src/test/resources/nexmark/q10_orc.sql new file mode 100644 index 00000000000..29df04b7340 --- /dev/null +++ b/gluten-flink/ut/src/test/resources/nexmark/q10_orc.sql @@ -0,0 +1,23 @@ +CREATE TABLE nexmark_q10_orc ( + auction BIGINT, + bidder BIGINT, + price BIGINT, + `dateTime` TIMESTAMP(3), + extra VARCHAR, + dt STRING, + hm STRING +) PARTITIONED BY (dt, hm) WITH ( + 'connector' = 'filesystem', + 'path' = 'file:///tmp/data/output/bid_orc/', + 'format' = 'orc', + 'sink.partition-commit.trigger' = 'process-time', + 'sink.partition-commit.delay' = '0s', + 'sink.partition-commit.policy.kind' = 'success-file', + 'partition.time-extractor.timestamp-pattern' = '$dt $hm:00', + 'sink.rolling-policy.rollover-interval' = '1s', + 'sink.rolling-policy.check-interval' = '1s' +); + +INSERT INTO nexmark_q10_orc +SELECT auction, bidder, price, `dateTime`, extra, DATE_FORMAT(`dateTime`, 'yyyy-MM-dd'), DATE_FORMAT(`dateTime`, 'HH:mm') +FROM bid;