Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 8 additions & 6 deletions .github/workflows/codecov.yml
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,13 @@ jobs:
runs-on: ubuntu-latest

steps:
- uses: actions/checkout@v1
- name: Set up JDK 1.8
uses: actions/setup-java@v1
- uses: actions/checkout@v4
- name: Set up JDK 11
uses: actions/setup-java@v4
with:
java-version: 1.8
distribution: temurin
java-version: 11
cache: maven
- name: Cache Maven Repository
uses: actions/cache@v3
with:
Expand All @@ -23,7 +25,7 @@ jobs:
- name: Build with Maven
run: ./mvnw test
- name: Codecov
uses: codecov/codecov-action@v2.1.0
uses: codecov/codecov-action@v4
with:
# Repository upload token - get it from codecov.io. Required only for private repositories
token: ${{ secrets.CODECOV_TOKEN }}
token: ${{ secrets.CODECOV_TOKEN }}
11 changes: 6 additions & 5 deletions .github/workflows/maven-publish.yml
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
name: Auto Deploy reactor-ql to the JetLinks Maven Repository
on:
push:
branches: ["master"]
branches: ["master","1.1"]

jobs:
publish:
Expand All @@ -12,13 +12,14 @@ jobs:
# os: [ubuntu-latest, windows-latest, macOS-latest]
os: [ ubuntu-latest ]
steps:
- uses: actions/checkout@v1
- uses: actions/checkout@v4
- run: echo ${{github.ref}}
- name: Set up Repository info
uses: actions/setup-java@v2
uses: actions/setup-java@v4
with:
java-version: '8'
distribution: 'adopt'
java-version: '11'
distribution: temurin
cache: maven
- name: Cache Maven Repository
uses: actions/cache@v3
with:
Expand Down
12 changes: 7 additions & 5 deletions .github/workflows/pull_request.yml
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,13 @@ jobs:
build:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v1
- name: Set up JDK 1.8
uses: actions/setup-java@v1
- uses: actions/checkout@v4
- name: Set up JDK 11
uses: actions/setup-java@v4
with:
java-version: 1.8
distribution: temurin
java-version: 11
cache: maven
- name: Cache Maven Repository
uses: actions/cache@v3
with:
Expand All @@ -23,4 +25,4 @@ jobs:
uses: codecov/codecov-action@v2.1.0
with:
# Repository upload token - get it from codecov.io. Required only for private repositories
token: ${{ secrets.CODECOV_TOKEN }}
token: ${{ secrets.CODECOV_TOKEN }}
7 changes: 7 additions & 0 deletions codecov.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
ignore:
- "pom.xml"
- ".github/**/*"
- "src/main/java/org/jetlinks/reactor/ql/supports/ExpressionVisitorAdapter.java"

github_checks:
annotations: false
14 changes: 7 additions & 7 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@

<groupId>org.jetlinks</groupId>
<artifactId>reactor-ql</artifactId>
<version>1.0.21-SNAPSHOT</version>
<version>1.1.0-SNAPSHOT</version>

<name>JetLinks</name>
<url>https://github.com/jetlinks/reactor-ql</url>
Expand Down Expand Up @@ -59,7 +59,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<project.build.locales>zh_CN</project.build.locales>
<java.version>1.8</java.version>
<java.version>11</java.version>
<project.build.jdk>${java.version}</project.build.jdk>
<reactor.version>2020.0.38</reactor.version>
</properties>
Expand Down Expand Up @@ -205,6 +205,7 @@
<excludes>
<exclude>**/ExpressionVisitorAdapter*</exclude>
<exclude>**/ReactorQLMetadata*</exclude>
<exclude>net/sf/jsqlparser/**</exclude>
</excludes>
</configuration>
<executions>
Expand Down Expand Up @@ -253,10 +254,9 @@
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.1</version>
<version>3.11.0</version>
<configuration>
<source>${project.build.jdk}</source>
<target>${project.build.jdk}</target>
<release>${project.build.jdk}</release>
<encoding>${project.build.sourceEncoding}</encoding>
</configuration>
</plugin>
Expand Down Expand Up @@ -288,7 +288,7 @@
<dependency>
<groupId>com.github.jsqlparser</groupId>
<artifactId>jsqlparser</artifactId>
<version>4.6</version>
<version>5.3</version>
</dependency>

<dependency>
Expand Down Expand Up @@ -397,4 +397,4 @@
<url>https://maven.aliyun.com/nexus/content/groups/public/</url>
</pluginRepository>
</pluginRepositories>
</project>
</project>
89 changes: 39 additions & 50 deletions src/main/java/org/jetlinks/reactor/ql/DefaultReactorQL.java
Original file line number Diff line number Diff line change
Expand Up @@ -185,12 +185,12 @@ protected Function<Flux<ReactorQLRecord>, Flux<ReactorQLRecord>> createJoin() {
Function<ReactorQLRecord, Flux<ReactorQLRecord>> rightStreamGetter = null;

//join (select deviceId,avg(temp) from temp group by interval('10s'),deviceId )
if (from instanceof SubSelect) {
if (from instanceof ParenthesedSelect) {
String alias = from.getAlias() == null ? null : from.getAlias().getName();
//子查询
DefaultReactorQL ql =
new DefaultReactorQL(new DefaultReactorQLMetadata(metadata,
((PlainSelect) ((SubSelect) from).getSelectBody())));
((PlainSelect) ((ParenthesedSelect) from).getSelect().getPlainSelect())));

rightStreamGetter = record -> ql
.builder
Expand Down Expand Up @@ -301,7 +301,8 @@ protected Function<Flux<ReactorQLRecord>, Flux<ReactorQLRecord>> createGroupBy()
groupByRef.set(nameMapper);
}
};
for (Expression groupByExpression : groupBy.getGroupByExpressionList().getExpressions()) {
for (Object expression : groupBy.getGroupByExpressionList().getExpressions()) {
Expression groupByExpression = (Expression) expression;
//函数分组, group by interval('1s')
if (groupByExpression instanceof net.sf.jsqlparser.expression.Function) {
featureConsumer.accept(null,
Expand Down Expand Up @@ -404,54 +405,42 @@ private Function<Flux<ReactorQLRecord>, Flux<ReactorQLRecord>> createMapper() {

List<Consumer<ReactorQLRecord>> allMapper = new ArrayList<>();

for (SelectItem selectItem : metadata.getSql().getSelectItems()) {
selectItem.accept(new SelectItemVisitorAdapter() {
// select a,b,c
@Override
public void visit(SelectExpressionItem item) {
Expression expression = item.getExpression();
String alias = item.getAlias() == null ? expression.toString() : item.getAlias().getName();
String fAlias = SqlUtils.getCleanStr(alias);
// select a,b,c
createExpressionMapper(expression).ifPresent(mapper -> mappers.put(fAlias, mapper));
// select count(),max(val)...
createAggMapper(expression).ifPresent(mapper -> aggMapper.put(fAlias, mapper));
//flatMap
ValueFlatMapFeature.createMapperByExpression(expression, metadata)
.ifPresent(mapper -> flatMappers.put(fAlias, mapper));

if (!mappers.containsKey(fAlias) && !aggMapper.containsKey(fAlias) && !flatMappers.containsKey(fAlias)) {
throw new UnsupportedOperationException("Unsupported expression:" + expression);
}
}

//select *
@Override
public void visit(AllColumns columns) {
allMapper.add(ReactorQLRecord::putRecordToResult);
}

//select t.*
@Override
public void visit(AllTableColumns columns) {
String name;
Alias alias = columns.getTable().getAlias();
if (alias == null) {
name = SqlUtils.getCleanStr(columns.getTable().getName());
} else {
name = SqlUtils.getCleanStr(alias.getName());
}
allMapper.add(record -> record
.getRecord(name)
.ifPresent(v -> {
if (v instanceof Map) {
record.setResults(((Map) v));
} else {
record.setResult(name, v);
}
}));
for (SelectItem<?> selectItem : metadata.getSql().getSelectItems()) {
Expression expression = selectItem.getExpression();
if (expression instanceof AllColumns) {
allMapper.add(ReactorQLRecord::putRecordToResult);
continue;
}
if (expression instanceof AllTableColumns) {
AllTableColumns columns = (AllTableColumns) expression;
String name;
Alias alias = columns.getTable().getAlias();
if (alias == null) {
name = SqlUtils.getCleanStr(columns.getTable().getName());
} else {
name = SqlUtils.getCleanStr(alias.getName());
}
});
allMapper.add(record -> record
.getRecord(name)
.ifPresent(v -> {
if (v instanceof Map) {
record.setResults(((Map) v));
} else {
record.setResult(name, v);
}
}));
continue;
}
String alias = selectItem.getAlias() == null ? expression.toString() : selectItem.getAlias().getName();
String fAlias = SqlUtils.getCleanStr(alias);
createExpressionMapper(expression).ifPresent(mapper -> mappers.put(fAlias, mapper));
createAggMapper(expression).ifPresent(mapper -> aggMapper.put(fAlias, mapper));
ValueFlatMapFeature.createMapperByExpression(expression, metadata)
.ifPresent(mapper -> flatMappers.put(fAlias, mapper));

if (!mappers.containsKey(fAlias) && !aggMapper.containsKey(fAlias) && !flatMappers.containsKey(fAlias)) {
throw new UnsupportedOperationException("Unsupported expression:" + expression);
}
}
Function<ReactorQLRecord, Mono<ReactorQLRecord>> _resultMapper;

Expand Down
72 changes: 31 additions & 41 deletions src/main/java/org/jetlinks/reactor/ql/feature/FromFeature.java
Original file line number Diff line number Diff line change
Expand Up @@ -40,56 +40,46 @@ static Function<ReactorQLContext, Flux<ReactorQLRecord>> createFromMapperByFrom(
if (body == null) {
return ctx -> ctx.getDataSource(null).map(val -> ReactorQLRecord.newRecord(null, val, ctx));
}
AtomicReference<Function<ReactorQLContext, Flux<ReactorQLRecord>>> ref = new AtomicReference<>();

body.accept(new FromItemVisitorAdapter() {
// from table
@Override
public void visit(Table table) {
ref.set(metadata.getFeatureNow(FeatureId.From.table)
.createFromMapper(table, metadata));
}

// from (select ...)
@Override
public void visit(SubSelect subSelect) {
ref.set(metadata.getFeatureNow(FeatureId.From.subSelect)
.createFromMapper(subSelect, metadata));
}

// select * from (values(6)) t(v)
@Override
public void visit(ValuesList valuesList) {
ref.set(metadata.getFeatureNow(FeatureId.From.values)
.createFromMapper(valuesList, metadata));
}

//select * from mysql(...)
@Override
public void visit(TableFunction tableFunction) {
ref.set(metadata
.getFeatureNow(FeatureId.From.of(tableFunction.getFunction().getName()),
tableFunction::toString)
.createFromMapper(tableFunction, metadata));
}

@Override
public void visit(ParenthesisFromItem aThis) {
ref.set(createFromMapperByFrom(aThis.getFromItem(), metadata));
if (body instanceof Table) {
return metadata.getFeatureNow(FeatureId.From.table)
.createFromMapper(body, metadata);
}
if (body instanceof ParenthesedSelect || body instanceof Select) {
return metadata.getFeatureNow(FeatureId.From.subSelect)
.createFromMapper(body, metadata);
}
if (body instanceof Values) {
return metadata.getFeatureNow(FeatureId.From.values)
.createFromMapper(body, metadata);
}
if (body instanceof ParenthesedFromItem) {
ParenthesedFromItem fromItem = (ParenthesedFromItem) body;
if (fromItem.getFromItem() instanceof Values) {
return metadata.getFeatureNow(FeatureId.From.values)
.createFromMapper(body, metadata);
}
});
if (ref.get() == null) {
throw new UnsupportedOperationException("不支持的查询:" + body);
return createFromMapperByFrom(fromItem.getFromItem(), metadata);
}
if (body instanceof TableFunction) {
TableFunction tableFunction = (TableFunction) body;
return metadata
.getFeatureNow(FeatureId.From.of(tableFunction.getFunction().getName()),
tableFunction::toString)
.createFromMapper(tableFunction, metadata);
}
return ref.get();
throw new UnsupportedOperationException("不支持的查询:" + body);
}

static Function<ReactorQLContext, Flux<ReactorQLRecord>> createFromMapperByBody(SelectBody body, ReactorQLMetadata metadata) {
static Function<ReactorQLContext, Flux<ReactorQLRecord>> createFromMapperByBody(Select body, ReactorQLMetadata metadata) {

FromItem from = null;
if (body instanceof PlainSelect) {
PlainSelect select = ((PlainSelect) body);
from = select.getFromItem();
} else if (body instanceof ParenthesedSelect) {
return createFromMapperByBody(((ParenthesedSelect) body).getSelect(), metadata);
} else if (body instanceof Values) {
from = body;
}
return createFromMapperByFrom(from, metadata);
}
Expand Down
15 changes: 12 additions & 3 deletions src/main/java/org/jetlinks/reactor/ql/feature/ValueMapFeature.java
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,9 @@

import net.sf.jsqlparser.expression.*;
import net.sf.jsqlparser.expression.operators.relational.ExistsExpression;
import net.sf.jsqlparser.expression.operators.relational.ExpressionList;
import net.sf.jsqlparser.schema.Column;
import net.sf.jsqlparser.statement.select.SubSelect;
import net.sf.jsqlparser.statement.select.Select;
import org.apache.commons.collections.CollectionUtils;
import org.jetlinks.reactor.ql.ReactorQLMetadata;
import org.jetlinks.reactor.ql.ReactorQLRecord;
Expand Down Expand Up @@ -77,7 +78,7 @@ public void visit(net.sf.jsqlparser.expression.Function function) {

//select (select * from xxx) data1 from ...
@Override
public void visit(SubSelect subSelect) {
public void visit(Select subSelect) {
ref.set(metadata
.getFeatureNow(FeatureId.ValueMap.select, expr::toString)
.createMapper(subSelect, metadata));
Expand Down Expand Up @@ -190,6 +191,12 @@ public void visit(DoubleValue value) {
ref.set((v) -> val);
}

@Override
public void visit(BooleanValue value) {
Mono<Object> val = Mono.just(value.getValue());
ref.set((v) -> val);
}

//select {d 'yyyy-mm-dd'}
@Override
public void visit(DateValue value) {
Expand Down Expand Up @@ -291,7 +298,7 @@ static Tuple2<Function<ReactorQLRecord, Publisher<?>>, Function<ReactorQLRecord,
List<Expression> expressions;
//只能有2个参数
if (function.getParameters() == null
|| CollectionUtils.isEmpty(expressions = function.getParameters().getExpressions())
|| CollectionUtils.isEmpty(expressions = ExpressionUtils.getFunctionParameter(function))
|| expressions.size() != 2) {
throw new IllegalArgumentException("The number of parameters must be 2 :" + expression);
}
Expand All @@ -301,6 +308,8 @@ static Tuple2<Function<ReactorQLRecord, Publisher<?>>, Function<ReactorQLRecord,
BinaryExpression bie = ((BinaryExpression) expression);
left = bie.getLeftExpression();
right = bie.getRightExpression();
} else if (expression instanceof ExpressionList && ((ExpressionList<?>) expression).size() == 1) {
return createBinaryMapper(((ExpressionList<Expression>) expression).get(0), metadata);
} else {
throw new UnsupportedOperationException("Unsupported expression:" + expression);
}
Expand Down
Loading