Skip to content

Commit 6bfd036

Browse files
committed
feat: Add RAW_JSON alias to the custom FlinkJsonType
1 parent 4581cd3 commit 6bfd036

9 files changed

Lines changed: 202 additions & 9 deletions

File tree

sqrl-planner/pom.xml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -149,6 +149,10 @@
149149
<groupId>com.datasqrl.flinkrunner</groupId>
150150
<artifactId>stdlib-commons</artifactId>
151151
</dependency>
152+
<dependency>
153+
<groupId>com.datasqrl.flinkrunner</groupId>
154+
<artifactId>stdlib-json</artifactId>
155+
</dependency>
152156
<dependency>
153157
<groupId>com.datasqrl.flinkrunner</groupId>
154158
<artifactId>stdlib-text</artifactId>

sqrl-planner/src/main/java/com/datasqrl/engine/stream/flink/FlinkSqlNodes.java

Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@
1616
package com.datasqrl.engine.stream.flink;
1717

1818
import com.datasqrl.calcite.schema.sql.SqlDataTypeSpecBuilder;
19+
import com.datasqrl.flinkrunner.stdlib.json.FlinkJsonType;
20+
import com.datasqrl.flinkrunner.stdlib.json.FlinkJsonTypeSerializer;
1921
import com.datasqrl.planner.util.NonSecretEnvVarResolver;
2022
import com.datasqrl.sql.SqlCallRewriter;
2123
import jakarta.annotation.Nullable;
@@ -33,6 +35,7 @@
3335
import org.apache.calcite.sql.SqlBasicCall;
3436
import org.apache.calcite.sql.SqlCall;
3537
import org.apache.calcite.sql.SqlCharStringLiteral;
38+
import org.apache.calcite.sql.SqlDataTypeSpec;
3639
import org.apache.calcite.sql.SqlIdentifier;
3740
import org.apache.calcite.sql.SqlIntervalQualifier;
3841
import org.apache.calcite.sql.SqlLiteral;
@@ -56,12 +59,18 @@
5659
import org.apache.flink.sql.parser.ddl.table.SqlTableLike;
5760
import org.apache.flink.sql.parser.ddl.view.SqlCreateView;
5861
import org.apache.flink.sql.parser.dml.RichSqlInsert;
62+
import org.apache.flink.sql.parser.type.SqlRawTypeNameSpec;
5963
import org.apache.flink.table.catalog.ObjectIdentifier;
64+
import org.apache.flink.table.types.logical.RawType;
6065

6166
public class FlinkSqlNodes {
6267

6368
public static final SqlDistribution NO_DISTRIBUTION = null;
6469

70+
private static final String RAW_JSON = "RAW_JSON";
71+
private static final RawType<FlinkJsonType> RAW_JSON_TYPE =
72+
new RawType<>(FlinkJsonType.class, new FlinkJsonTypeSerializer());
73+
6574
public static SqlIdentifier identifier(String str) {
6675
return new SqlIdentifier(str, SqlParserPos.ZERO);
6776
}
@@ -191,6 +200,109 @@ public static SqlCreateTable resolveTableProperties(SqlCreateTable createTable)
191200
createTable.ifNotExists);
192201
}
193202

203+
/**
204+
* Replaces the RAW_JSON column type alias with the actual RAW type used by the Flink JSON
205+
* functions.
206+
*/
207+
public static SqlCreateTable resolveRawJsonTypAliases(SqlCreateTable createTable) {
208+
var columns = new ArrayList<SqlNode>(createTable.getColumnList().size());
209+
var changed = false;
210+
for (var column : createTable.getColumnList()) {
211+
var resolvedColumn = resolveRawJsonType(column);
212+
columns.add(resolvedColumn);
213+
changed |= resolvedColumn != column;
214+
}
215+
216+
if (!changed) {
217+
return createTable;
218+
}
219+
220+
var resolvedColumns = new SqlNodeList(columns, createTable.getColumnList().getParserPosition());
221+
if (createTable instanceof SqlCreateTableLike likeTable) {
222+
return new SqlCreateTableLike(
223+
likeTable.getParserPosition(),
224+
likeTable.getName(),
225+
resolvedColumns,
226+
likeTable.getTableConstraints(),
227+
createProperties(likeTable.getProperties()),
228+
likeTable.getDistribution(),
229+
createPartitionKeys(likeTable.getPartitionKeyList()),
230+
likeTable.getWatermark().orElse(null),
231+
createStringLiteral(likeTable.getComment()),
232+
likeTable.getTableLike(),
233+
likeTable.isTemporary(),
234+
likeTable.ifNotExists);
235+
}
236+
237+
return new SqlCreateTable(
238+
createTable.getParserPosition(),
239+
createTable.getName(),
240+
resolvedColumns,
241+
createTable.getTableConstraints(),
242+
createProperties(createTable.getProperties()),
243+
createTable.getDistribution(),
244+
createPartitionKeys(createTable.getPartitionKeyList()),
245+
createTable.getWatermark().orElse(null),
246+
createStringLiteral(createTable.getComment()),
247+
createTable.isTemporary(),
248+
createTable.ifNotExists);
249+
}
250+
251+
private static SqlNode resolveRawJsonType(SqlNode column) {
252+
if (column instanceof SqlRegularColumn regularColumn
253+
&& isRawJsonType(regularColumn.getType())) {
254+
255+
return new SqlRegularColumn(
256+
regularColumn.getParserPosition(),
257+
regularColumn.getName(),
258+
createStringLiteral(regularColumn.getComment()),
259+
createFlexibleJsonRawType(regularColumn.getType()),
260+
regularColumn.getConstraint().orElse(null));
261+
}
262+
263+
if (column instanceof SqlMetadataColumn metadataColumn
264+
&& isRawJsonType(metadataColumn.getType())) {
265+
266+
var metadataAlias =
267+
metadataColumn.getMetadataAlias().map(FlinkSqlNodes::createStringLiteral).orElse(null);
268+
269+
return new SqlMetadataColumn(
270+
metadataColumn.getParserPosition(),
271+
metadataColumn.getName(),
272+
createStringLiteral(metadataColumn.getComment()),
273+
createFlexibleJsonRawType(metadataColumn.getType()),
274+
metadataAlias,
275+
metadataColumn.isVirtual());
276+
}
277+
278+
return column;
279+
}
280+
281+
private static boolean isRawJsonType(SqlDataTypeSpec type) {
282+
var typeName = type.getTypeNameSpec().getTypeName();
283+
return typeName != null && RAW_JSON.equalsIgnoreCase(typeName.getSimple());
284+
}
285+
286+
private static SqlDataTypeSpec createFlexibleJsonRawType(SqlDataTypeSpec originalType) {
287+
var originalPosition = originalType.getParserPosition();
288+
var rawTypeSql = RAW_JSON_TYPE.asSerializableString();
289+
290+
var position =
291+
new SqlParserPos(
292+
originalPosition.getLineNum(),
293+
originalPosition.getColumnNum(),
294+
originalPosition.getLineNum(),
295+
originalPosition.getColumnNum() + rawTypeSql.length() - 1);
296+
297+
var rawTypeName =
298+
new SqlRawTypeNameSpec(
299+
SqlLiteral.createCharString(RAW_JSON_TYPE.getOriginatingClass().getName(), position),
300+
SqlLiteral.createCharString(RAW_JSON_TYPE.getSerializerString(), position),
301+
position);
302+
303+
return new SqlDataTypeSpec(rawTypeName, position).withNullable(originalType.getNullable());
304+
}
305+
194306
public static SqlNodeList createProperties(Map<String, String> options) {
195307
var sqlNodes = new ArrayList<SqlNode>(options.size());
196308

sqrl-planner/src/main/java/com/datasqrl/planner/Sqrl2FlinkSQLTranslator.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -837,6 +837,7 @@ private AddTableResult addTable(
837837
var tableSqlNode = parseSQL(createTableSql);
838838
checkArgument(tableSqlNode instanceof SqlCreateTable, "Expected CREATE TABLE statement");
839839
var tableDefinition = FlinkSqlNodes.resolveTableProperties((SqlCreateTable) tableSqlNode);
840+
tableDefinition = FlinkSqlNodes.resolveRawJsonTypAliases(tableDefinition);
840841
var fullTable = tableDefinition;
841842
var origTableName = fullTable.getName().getSimple();
842843
final var finalTableName = tableNameModifier.apply(origTableName);

sqrl-planner/src/test/java/com/datasqrl/engine/stream/flink/plan/FlinkSqlNodesTest.java

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,11 +26,19 @@
2626
import java.util.List;
2727
import java.util.Map;
2828
import java.util.Optional;
29+
import org.apache.calcite.sql.SqlDataTypeSpec;
2930
import org.apache.calcite.sql.SqlIdentifier;
31+
import org.apache.calcite.sql.SqlLiteral;
3032
import org.apache.calcite.sql.SqlNode;
3133
import org.apache.calcite.sql.SqlNodeList;
3234
import org.apache.calcite.sql.SqlSelect;
35+
import org.apache.calcite.sql.SqlUserDefinedTypeNameSpec;
3336
import org.apache.calcite.sql.parser.SqlParserPos;
37+
import org.apache.flink.sql.parser.ddl.SqlTableColumn.SqlMetadataColumn;
38+
import org.apache.flink.sql.parser.ddl.SqlTableColumn.SqlRegularColumn;
39+
import org.apache.flink.sql.parser.ddl.table.SqlCreateTableLike;
40+
import org.apache.flink.sql.parser.ddl.table.SqlTableLike;
41+
import org.apache.flink.sql.parser.type.SqlRawTypeNameSpec;
3442
import org.junit.jupiter.api.Test;
3543

3644
class FlinkSqlNodesTest {
@@ -174,4 +182,69 @@ void createPartitionKeys() {
174182
var expectedSql = "`year`, `month`, `day`";
175183
assertThat(sql.trim()).isEqualTo(expectedSql);
176184
}
185+
186+
@Test
187+
void resolveRawJsonTypAliases() {
188+
var position = new SqlParserPos(3, 20, 3, 32);
189+
var rawJsonType =
190+
new SqlDataTypeSpec(new SqlUserDefinedTypeNameSpec("RAW_JSON", position), position)
191+
.withNullable(false);
192+
var regularColumn =
193+
new SqlRegularColumn(
194+
position,
195+
FlinkSqlNodes.identifier("payload"),
196+
FlinkSqlNodes.createStringLiteral("payload comment"),
197+
rawJsonType,
198+
null);
199+
var metadataColumn =
200+
new SqlMetadataColumn(
201+
position,
202+
FlinkSqlNodes.identifier("event_time"),
203+
FlinkSqlNodes.createStringLiteral("metadata comment"),
204+
rawJsonType,
205+
SqlLiteral.createCharString("timestamp", position),
206+
true);
207+
var tableLike = new SqlTableLike(position, FlinkSqlNodes.identifier("base_table"), List.of());
208+
var table =
209+
new SqlCreateTableLike(
210+
position,
211+
FlinkSqlNodes.identifier("source_table"),
212+
new SqlNodeList(List.of(regularColumn, metadataColumn), position),
213+
List.of(),
214+
FlinkSqlNodes.createProperties(Map.of("connector", "kafka")),
215+
FlinkSqlNodes.NO_DISTRIBUTION,
216+
SqlNodeList.EMPTY,
217+
null,
218+
FlinkSqlNodes.createStringLiteral("table comment"),
219+
tableLike,
220+
true,
221+
true);
222+
223+
var resolved = FlinkSqlNodes.resolveRawJsonTypAliases(table);
224+
225+
assertThat(resolved).isInstanceOf(SqlCreateTableLike.class);
226+
assertThat(((SqlCreateTableLike) resolved).getTableLike()).isSameAs(tableLike);
227+
assertThat(resolved.getProperties()).isEqualTo(table.getProperties());
228+
assertThat(resolved.getColumnList().get(0)).isInstanceOf(SqlRegularColumn.class);
229+
assertThat(resolved.getColumnList().get(1)).isInstanceOf(SqlMetadataColumn.class);
230+
231+
var resolvedRegularColumn = (SqlRegularColumn) resolved.getColumnList().get(0);
232+
assertThat(resolvedRegularColumn.getType().getTypeNameSpec())
233+
.isInstanceOf(SqlRawTypeNameSpec.class);
234+
assertThat(resolvedRegularColumn.getType().getNullable()).isFalse();
235+
assertThat(resolvedRegularColumn.getComment()).isEqualTo("payload comment");
236+
assertThat(resolvedRegularColumn.getType().getParserPosition().getLineNum())
237+
.isEqualTo(position.getLineNum());
238+
assertThat(resolvedRegularColumn.getType().getParserPosition().getColumnNum())
239+
.isEqualTo(position.getColumnNum());
240+
assertThat(resolvedRegularColumn.getType().getParserPosition().getEndColumnNum())
241+
.isEqualTo(position.getColumnNum() + unparse(resolvedRegularColumn.getType()).length() - 1);
242+
243+
var resolvedMetadataColumn = (SqlMetadataColumn) resolved.getColumnList().get(1);
244+
assertThat(resolvedMetadataColumn.getType().getTypeNameSpec())
245+
.isInstanceOf(SqlRawTypeNameSpec.class);
246+
assertThat(resolvedMetadataColumn.getMetadataAlias()).contains("timestamp");
247+
assertThat(resolvedMetadataColumn.isVirtual()).isTrue();
248+
assertThat(resolvedMetadataColumn.getComment()).isEqualTo("metadata comment");
249+
}
177250
}

sqrl-testing/sqrl-testing-integration/src/test/java/com/datasqrl/FullUseCaseIT.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,8 +65,7 @@ void specificUseCase(UseCaseParam param, TestContainerHook hook) {
6565

6666
/** Ad-hoc debugging entry point. Change the path below to run a single use case manually. */
6767
static Stream<UseCaseParam> specificUseCaseProvider() {
68-
return Stream.of(
69-
new UseCaseParam(USE_CASES.resolve("function-translation/duckdb").resolve("package.json")));
68+
return Stream.of(new UseCaseParam(USE_CASES.resolve("jwt-authorized").resolve("package.json")));
7069
}
7170

7271
@ParameterizedTest

sqrl-testing/sqrl-testing-integration/src/test/resources/usecases/jwt-authorized/jwt-authorized.sqrl

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,11 @@ IMPORT jwt-authorized-base.*;
44
CREATE TABLE AuthInputData (
55
event_id STRING NOT NULL METADATA FROM 'uuid',
66
val BIGINT NOT NULL METADATA FROM 'auth.val',
7-
message STRING NOT NULL,
7+
payload RAW_JSON,
88
event_time TIMESTAMP_LTZ(3) NOT NULL METADATA FROM 'timestamp'
99
);
1010

11-
_Messages := SELECT * FROM AuthInputData WHERE message <> '';
11+
_Messages := SELECT * FROM AuthInputData WHERE payload IS NOT NULL;
1212

1313
MessageSubscription(val BIGINT NOT NULL METADATA FROM 'auth.val') :=
1414
SUBSCRIBE SELECT * FROM _Messages WHERE val = :val;

sqrl-testing/sqrl-testing-integration/src/test/resources/usecases/jwt-authorized/snapshots/AuthInputData-valid-subscription.snapshot

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,14 +2,18 @@
22
"data" : {
33
"MessageSubscription" : {
44
"val" : 1,
5-
"message" : "Hello"
5+
"payload" : {
6+
"message" : "Hello"
7+
}
68
}
79
}
810
}, {
911
"data" : {
1012
"MessageSubscription" : {
1113
"val" : 1,
12-
"message" : "World"
14+
"payload" : {
15+
"message" : "World"
16+
}
1317
}
1418
}
1519
} ]
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
mutation {
2-
mut1: AuthInputData(event: {message: "Hello"}) {
2+
mut1: AuthInputData(event: {payload: {message: "Hello"}}) {
33
val
44
}
5-
mut2: AuthInputData(event: {message: "World"}) {
5+
mut2: AuthInputData(event: {payload: {message: "World"}}) {
66
val
77
}
88
}
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
subscription {
22
MessageSubscription {
33
val
4-
message
4+
payload
55
}
66
}

0 commit comments

Comments
 (0)