Skip to content

Commit e4226a0

Browse files
committed
Support ksqlDB window grace periods and EMIT FINAL
1 parent b2115ac commit e4226a0

7 files changed

Lines changed: 360 additions & 70 deletions

File tree

README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,7 @@ and missing syntax gets added on demand — [open an issue](https://github.com/J
141141
| **DML** | `INSERT` · `UPDATE` · `UPSERT` · `MERGE` · `DELETE` · `TRUNCATE TABLE` |
142142
| **DDL** | `CREATE …` · `ALTER …` · `DROP …` |
143143
| **PostgreSQL RLS** | `CREATE POLICY` · `ALTER TABLE … ENABLE`/`DISABLE`/`FORCE`/`NO FORCE ROW LEVEL SECURITY` |
144+
| **ksqlDB windows** | JOIN `WITHIN`, window `GRACE PERIOD`, and `EMIT CHANGES`/`FINAL` |
144145
| **Salesforce SOQL** | `INCLUDES` · `EXCLUDES` |
145146

146147
Beyond statement shapes, the grammar handles nested sub-selects, bind parameters (`?`,

src/main/java/net/sf/jsqlparser/statement/select/KSQLJoinWindow.java

Lines changed: 44 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,6 @@
99
*/
1010
package net.sf.jsqlparser.statement.select;
1111

12-
import java.util.Locale;
1312
import net.sf.jsqlparser.parser.ASTNodeAccessImpl;
1413

1514
import static net.sf.jsqlparser.statement.select.KSQLWindow.TimeUnit;
@@ -23,9 +22,37 @@ public class KSQLJoinWindow extends ASTNodeAccessImpl {
2322
private TimeUnit beforeTimeUnit;
2423
private long afterDuration;
2524
private TimeUnit afterTimeUnit;
25+
private boolean usingBrackets = true;
26+
private KSQLWindow.Duration gracePeriod;
27+
28+
public boolean isUsingBrackets() {
29+
return usingBrackets || beforeAfter;
30+
}
31+
32+
public void setUsingBrackets(boolean usingBrackets) {
33+
this.usingBrackets = usingBrackets;
34+
}
35+
36+
public KSQLJoinWindow withUsingBrackets(boolean usingBrackets) {
37+
setUsingBrackets(usingBrackets);
38+
return this;
39+
}
40+
41+
public KSQLWindow.Duration getGracePeriod() {
42+
return gracePeriod;
43+
}
44+
45+
public void setGracePeriod(KSQLWindow.Duration gracePeriod) {
46+
this.gracePeriod = gracePeriod;
47+
}
48+
49+
public KSQLJoinWindow withGracePeriod(KSQLWindow.Duration gracePeriod) {
50+
setGracePeriod(gracePeriod);
51+
return this;
52+
}
2653

2754
public final static TimeUnit from(String timeUnitStr) {
28-
return Enum.valueOf(TimeUnit.class, timeUnitStr.toUpperCase(Locale.ROOT));
55+
return TimeUnit.from(timeUnitStr);
2956
}
3057

3158
public boolean isBeforeAfterWindow() {
@@ -86,11 +113,23 @@ public void setAfterTimeUnit(TimeUnit afterTimeUnit) {
86113

87114
@Override
88115
public String toString() {
116+
StringBuilder builder = new StringBuilder();
117+
if (isUsingBrackets()) {
118+
builder.append('(');
119+
}
89120
if (isBeforeAfterWindow()) {
90-
return "(" + beforeDuration + " " + beforeTimeUnit + ", " + afterDuration + " "
91-
+ afterTimeUnit + ")";
121+
builder.append(beforeDuration).append(' ').append(beforeTimeUnit)
122+
.append(", ").append(afterDuration).append(' ').append(afterTimeUnit);
123+
} else {
124+
builder.append(duration).append(' ').append(timeUnit);
125+
}
126+
if (isUsingBrackets()) {
127+
builder.append(')');
128+
}
129+
if (gracePeriod != null) {
130+
builder.append(" GRACE PERIOD ").append(gracePeriod);
92131
}
93-
return "(" + duration + " " + timeUnit + ")";
132+
return builder.toString();
94133
}
95134

96135
public KSQLJoinWindow withDuration(long duration) {

src/main/java/net/sf/jsqlparser/statement/select/KSQLWindow.java

Lines changed: 51 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,47 @@ public class KSQLWindow extends ASTNodeAccessImpl {
2121
private TimeUnit sizeTimeUnit;
2222
private long advanceDuration;
2323
private TimeUnit advanceTimeUnit;
24+
private Duration gracePeriod;
25+
26+
/** A non-negative duration with its SQL time unit, including an explicit zero. */
27+
public static final class Duration implements java.io.Serializable {
28+
private final long value;
29+
private final TimeUnit timeUnit;
30+
31+
public Duration(long value, TimeUnit timeUnit) {
32+
if (value < 0) {
33+
throw new IllegalArgumentException("Duration must not be negative");
34+
}
35+
this.value = value;
36+
this.timeUnit = java.util.Objects.requireNonNull(timeUnit, "timeUnit");
37+
}
38+
39+
public long getValue() {
40+
return value;
41+
}
42+
43+
public TimeUnit getTimeUnit() {
44+
return timeUnit;
45+
}
46+
47+
@Override
48+
public String toString() {
49+
return value + " " + timeUnit;
50+
}
51+
}
52+
53+
public Duration getGracePeriod() {
54+
return gracePeriod;
55+
}
56+
57+
public void setGracePeriod(Duration gracePeriod) {
58+
this.gracePeriod = gracePeriod;
59+
}
60+
61+
public KSQLWindow withGracePeriod(Duration gracePeriod) {
62+
setGracePeriod(gracePeriod);
63+
return this;
64+
}
2465

2566
public KSQLWindow() {}
2667

@@ -82,14 +123,20 @@ public void setAdvanceTimeUnit(TimeUnit advanceTimeUnit) {
82123

83124
@Override
84125
public String toString() {
126+
StringBuilder builder = new StringBuilder();
85127
if (isHoppingWindow()) {
86-
return "HOPPING (" + "SIZE " + sizeDuration + " " + sizeTimeUnit + ", " +
87-
"ADVANCE BY " + advanceDuration + " " + advanceTimeUnit + ")";
128+
builder.append("HOPPING (SIZE ").append(sizeDuration).append(' ').append(sizeTimeUnit)
129+
.append(", ADVANCE BY ").append(advanceDuration).append(' ')
130+
.append(advanceTimeUnit);
88131
} else if (isSessionWindow()) {
89-
return "SESSION (" + sizeDuration + " " + sizeTimeUnit + ")";
132+
builder.append("SESSION (").append(sizeDuration).append(' ').append(sizeTimeUnit);
90133
} else {
91-
return "TUMBLING (" + "SIZE " + sizeDuration + " " + sizeTimeUnit + ")";
134+
builder.append("TUMBLING (SIZE ").append(sizeDuration).append(' ').append(sizeTimeUnit);
135+
}
136+
if (gracePeriod != null) {
137+
builder.append(", GRACE PERIOD ").append(gracePeriod);
92138
}
139+
return builder.append(')').toString();
93140
}
94141

95142
public KSQLWindow withSizeDuration(long sizeDuration) {

src/main/java/net/sf/jsqlparser/statement/select/PlainSelect.java

Lines changed: 29 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,8 @@ public class PlainSelect extends Select {
5555
private MySqlSqlCacheFlags mySqlCacheFlag = null;
5656
private String forXmlPath;
5757
private KSQLWindow ksqlWindow = null;
58-
private boolean emitChanges = false;
58+
private EmitMode emitMode = EmitMode.NONE;
59+
5960
private List<WindowDefinition> windowDefinitions;
6061
/**
6162
* @see <a href=
@@ -68,6 +69,10 @@ public class PlainSelect extends Select {
6869
private Table intoTempTable = null;
6970
private List<UpdateSet> settings = null;
7071

72+
public enum EmitMode {
73+
NONE, CHANGES, FINAL
74+
}
75+
7176
public PlainSelect() {}
7277

7378
public PlainSelect(FromItem fromItem) {
@@ -500,12 +505,32 @@ public void setKsqlWindow(KSQLWindow ksqlWindow) {
500505
this.ksqlWindow = ksqlWindow;
501506
}
502507

508+
public EmitMode getEmitMode() {
509+
return emitMode;
510+
}
511+
512+
public void setEmitMode(EmitMode emitMode) {
513+
this.emitMode = java.util.Objects.requireNonNull(emitMode, "emitMode");
514+
}
515+
516+
public PlainSelect withEmitMode(EmitMode emitMode) {
517+
setEmitMode(emitMode);
518+
return this;
519+
}
520+
521+
public StringBuilder appendEmitClauseTo(StringBuilder builder) {
522+
if (emitMode != EmitMode.NONE) {
523+
builder.append(" EMIT ").append(emitMode);
524+
}
525+
return builder;
526+
}
527+
503528
public boolean isEmitChanges() {
504-
return emitChanges;
529+
return emitMode == EmitMode.CHANGES;
505530
}
506531

507532
public void setEmitChanges(boolean emitChanges) {
508-
this.emitChanges = emitChanges;
533+
emitMode = emitChanges ? EmitMode.CHANGES : EmitMode.NONE;
509534
}
510535

511536
public List<WindowDefinition> getWindowDefinitions() {
@@ -634,9 +659,7 @@ public StringBuilder appendSelectBodyTo(StringBuilder builder) {
634659
builder.append(windowDefinitions.stream().map(WindowDefinition::toString)
635660
.collect(joining(", ")));
636661
}
637-
if (emitChanges) {
638-
builder.append(" EMIT CHANGES");
639-
}
662+
appendEmitClauseTo(builder);
640663
if (intoTempTable != null) {
641664
builder.append(" INTO TEMP ").append(intoTempTable);
642665
}

src/main/java/net/sf/jsqlparser/util/deparser/SelectDeParser.java

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -354,9 +354,7 @@ public <S> StringBuilder visit(PlainSelect plainSelect, S context) {
354354
builder.append(plainSelect.getOption());
355355
}
356356

357-
if (plainSelect.isEmitChanges()) {
358-
builder.append(" EMIT CHANGES");
359-
}
357+
plainSelect.appendEmitClauseTo(builder);
360358
if (plainSelect.getLimitBy() != null) {
361359
new LimitDeparser(expressionVisitor, builder).deParse(plainSelect.getLimitBy());
362360
}

0 commit comments

Comments
 (0)