fqaiser94 commented on code in PR #174:
URL: 
https://github.com/apache/flink-connector-kafka/pull/174#discussion_r4123696471


##########
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaTableITCase.java:
##########
@@ -1721,6 +1750,504 @@ private void 
testStartFromGroupOffsetsWithNoneResetStrategy(final String format)
         }
     }
 
+    private void projectionPushdownSetupData(final String format, final String 
topic)
+            throws Exception {
+        createTestTopic(topic, 1, 1);
+
+        String groupId = getStandardProps().getProperty("group.id");
+        String bootstraps = getBootstrapServers();
+
+        final String createTable =
+                String.format(
+                        "CREATE TABLE kafka (\n"
+                                + "  `a` STRING,\n"
+                                + "  `b` STRING,\n"
+                                + "  `topic` STRING NOT NULL METADATA 
VIRTUAL,\n"
+                                + "  `c` STRING,\n"
+                                + "  `partition` INT NOT NULL METADATA 
VIRTUAL,\n"
+                                + "  `d` STRING\n"
+                                + ") WITH (\n"
+                                + "  'connector' = 'kafka',\n"
+                                + "  'topic' = '%s',\n"
+                                + "  'properties.bootstrap.servers' = '%s',\n"
+                                + "  'properties.group.id' = '%s',\n"
+                                + "  'scan.startup.mode' = 
'earliest-offset',\n"
+                                + "  %s,\n"
+                                + "  'key.fields' = 'a; b',\n"
+                                + "  %s,\n"
+                                + "  'value.fields-include' = 'EXCEPT_KEY'\n"
+                                + ")",
+                        topic,
+                        bootstraps,
+                        groupId,
+                        keyFormatOptions(format),
+                        valueFormatOptions(format));
+        tEnv.executeSql(createTable);
+
+        final String initialValues = "INSERT INTO kafka (a, b, c, d) SELECT 
'a', 'b', 'c', 'd'";
+        tEnv.executeSql(initialValues).await();
+    }
+
+    @ParameterizedTest(name = "format: {0}")
+    @MethodSource("formats")
+    public void testProjectionPushdownSelectAllFields(final String format) 
throws Exception {
+        final String topic = "testProjectionPushdown_" + format + "_" + 
UUID.randomUUID();
+        projectionPushdownSetupData(format, topic);
+
+        assertQueryResult(
+                "SELECT * FROM kafka",
+                "== Optimized Execution Plan ==\n"
+                        + "Calc(select=[a, b, topic, c, partition, d])\n"
+                        + "+- TableSourceScan(table=[[default_catalog, 
default_database, kafka, metadata=[topic, partition]]], fields=[a, b, c, d, 
topic, partition])\n",
+                Collections.singletonList(String.format("+I(a,b,%s,c,%d,d)", 
topic, 0)));
+
+        cleanupTopic(topic);
+    }
+
+    @ParameterizedTest(name = "format: {0}")
+    @MethodSource("formats")
+    public void testProjectionPushdownSelectSpecificPhysicalFields(final 
String format)

Review Comment:
   Fixed. 
   
   Added 2 new tests to cover these cases: 
   - `testProjectionPushdownSelectOnlyValueFields` in 
https://github.com/apache/flink-connector-kafka/pull/174/commits/619b693747456502b25589b1e1eff1ea0a472092
   - `testProjectionPushdownSkipsRecordsWithNullKey` in 
https://github.com/apache/flink-connector-kafka/pull/174/commits/692840dcf6953974a420fb66eced473b428b10ce



-- 
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