Skip to content

Commit 92e9f65

Browse files
committed
feat(timescaledb): 支持配置schema并优化时间分组查询
- 在TimescaleDBOperations接口及实现类中添加schema()方法 - 修改TimescaleDBColumnModeQueryOperations和RowModeQueryOperations构造函数,增加schema参数 - 更新doAggregation0方法签名,传入schema参数用于时间分组查询 - 调整TimescaleDBUtils.createTimeGroupColumn方法,支持指定schema前缀 -优化TimescaleDBTimeSeriesService中的时间分组查询逻辑 - 调整类导入顺序,移除未使用的导入包
1 parent 70e42b1 commit 92e9f65

8 files changed

Lines changed: 43 additions & 28 deletions

File tree

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/TimescaleDBOperations.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@ public interface TimescaleDBOperations {
2121

2222
DatabaseOperator database();
2323

24+
String schema();
25+
2426
TimescaleDBDataWriter writer();
2527

2628
}

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/TimescaleDBUtils.java

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -20,17 +20,16 @@
2020
import org.hswebframework.ezorm.rdb.operator.dml.query.NativeSelectColumn;
2121
import org.hswebframework.ezorm.rdb.operator.dml.query.SelectColumn;
2222
import org.hswebframework.ezorm.rdb.supports.postgres.JsonbType;
23-
import org.jetlinks.core.metadata.DataType;
24-
import org.jetlinks.core.metadata.PropertyMetadata;
25-
import org.jetlinks.core.metadata.types.ArrayType;
26-
import org.jetlinks.core.metadata.types.ObjectType;
2723
import org.jetlinks.community.Interval;
2824
import org.jetlinks.community.things.data.ThingsDataConstants;
2925
import org.jetlinks.community.things.utils.ThingsDatabaseUtils;
3026
import org.jetlinks.community.timescaledb.metadata.JsonbValueCodec;
3127
import org.jetlinks.community.timeseries.TimeSeriesData;
3228
import org.jetlinks.community.timeseries.query.Aggregation;
33-
import org.jetlinks.community.utils.ObjectMappers;
29+
import org.jetlinks.core.metadata.DataType;
30+
import org.jetlinks.core.metadata.PropertyMetadata;
31+
import org.jetlinks.core.metadata.types.ArrayType;
32+
import org.jetlinks.core.metadata.types.ObjectType;
3433
import org.jetlinks.reactor.ql.utils.CastUtils;
3534

3635
public class TimescaleDBUtils {
@@ -40,15 +39,15 @@ public static String getTableName(String name) {
4039
return ThingsDatabaseUtils.createTableName(name);
4140
}
4241

43-
public static NativeSelectColumn createTimeGroupColumn(long startWith, Interval interval) {
42+
public static NativeSelectColumn createTimeGroupColumn(long startWith, Interval interval, String schema) {
4443

4544
String unit = interval.getNumber().intValue() + " " + interval
4645
.getUnit()
4746
.name()
4847
.toLowerCase();
4948

5049
return NativeSelectColumn
51-
.of("time_bucket('" + unit + "',timestamp)");
50+
.of(schema + "." + "time_bucket('" + unit + "',timestamp)");
5251
}
5352

5453
public static TimeSeriesData convertToTimeSeriesData(Record record) {

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/impl/DefaultTimescaleDBOperations.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,11 @@ public DatabaseOperator database() {
102102
return database;
103103
}
104104

105+
@Override
106+
public String schema() {
107+
return properties.getSchema();
108+
}
109+
105110
@Override
106111
public TimescaleDBDataWriter writer() {
107112
return writer;

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/thing/TimescaleDBColumnModeQueryOperations.java

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -28,26 +28,27 @@
2828
import org.hswebframework.web.api.crud.entity.PagerResult;
2929
import org.hswebframework.web.api.crud.entity.QueryParamEntity;
3030
import org.hswebframework.web.crud.query.QueryHelper;
31-
import org.jetlinks.core.metadata.EventMetadata;
32-
import org.jetlinks.core.things.ThingsRegistry;
3331
import org.jetlinks.community.things.data.AggregationRequest;
3432
import org.jetlinks.community.things.data.PropertyAggregation;
3533
import org.jetlinks.community.things.data.ThingsDataConstants;
3634
import org.jetlinks.community.things.data.ThingsDataUtils;
3735
import org.jetlinks.community.things.data.operations.ColumnModeQueryOperationsBase;
3836
import org.jetlinks.community.things.data.operations.DataSettings;
3937
import org.jetlinks.community.things.data.operations.MetricBuilder;
40-
import org.jetlinks.community.things.data.operations.RowModeQueryOperationsBase;
4138
import org.jetlinks.community.timescaledb.TimescaleDBUtils;
4239
import org.jetlinks.community.timeseries.TimeSeriesData;
4340
import org.jetlinks.community.timeseries.query.Aggregation;
4441
import org.jetlinks.community.timeseries.query.AggregationData;
42+
import org.jetlinks.core.things.ThingsRegistry;
4543
import org.jetlinks.reactor.ql.utils.CastUtils;
4644
import org.slf4j.Logger;
4745
import reactor.core.publisher.Flux;
4846
import reactor.core.publisher.Mono;
4947

50-
import java.util.*;
48+
import java.util.ArrayList;
49+
import java.util.Map;
50+
import java.util.NavigableMap;
51+
import java.util.Objects;
5152
import java.util.function.Function;
5253

5354
import static org.jetlinks.community.timescaledb.TimescaleDBUtils.createTimeGroupColumn;
@@ -56,15 +57,19 @@
5657
public class TimescaleDBColumnModeQueryOperations extends ColumnModeQueryOperationsBase {
5758
private final DatabaseOperator database;
5859

60+
private final String schema;
61+
5962
public TimescaleDBColumnModeQueryOperations(String thingType,
6063
String thingTemplateId,
6164
String thingId,
6265
MetricBuilder metricBuilder,
6366
DataSettings settings,
6467
ThingsRegistry registry,
65-
DatabaseOperator database) {
68+
DatabaseOperator database,
69+
String schema) {
6670
super(thingType, thingTemplateId, thingId, metricBuilder, settings, registry);
6771
this.database = database;
72+
this.schema = schema;
6873
}
6974

7075
@Override
@@ -104,12 +109,13 @@ record -> mapper.apply(convertToTimeSeriesData(record))
104109
protected Flux<AggregationData> doAggregation(String metric,
105110
AggregationRequest request,
106111
AggregationContext context) {
107-
return doAggregation0(database, metric, request, context);
112+
return doAggregation0(database, metric, schema, request, context);
108113
}
109114

110115

111116
static Flux<AggregationData> doAggregation0(DatabaseOperator database,
112117
String metric,
118+
String schema,
113119
AggregationRequest request,
114120
AggregationContext context) {
115121
metric = TimescaleDBUtils.getTableName(metric);
@@ -122,7 +128,8 @@ static Flux<AggregationData> doAggregation0(DatabaseOperator database,
122128
if (request.getInterval() != null) {
123129
NativeSelectColumn column = createTimeGroupColumn(
124130
request.getFrom().getTime(),
125-
request.getInterval()
131+
request.getInterval(),
132+
schema
126133
);
127134
query.groupBy(column);
128135
query.select(column);

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/thing/TimescaleDBColumnModeStrategy.java

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,13 @@
1515
*/
1616
package org.jetlinks.community.timescaledb.thing;
1717

18-
import org.jetlinks.community.things.data.TableSafeMetricBuilder;
19-
import org.jetlinks.core.things.ThingsRegistry;
2018
import org.jetlinks.community.things.data.AbstractThingDataRepositoryStrategy;
19+
import org.jetlinks.community.things.data.TableSafeMetricBuilder;
2120
import org.jetlinks.community.things.data.operations.DDLOperations;
2221
import org.jetlinks.community.things.data.operations.QueryOperations;
2322
import org.jetlinks.community.things.data.operations.SaveOperations;
2423
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
24+
import org.jetlinks.core.things.ThingsRegistry;
2525

2626
public class TimescaleDBColumnModeStrategy extends AbstractThingDataRepositoryStrategy {
2727

@@ -70,7 +70,8 @@ protected QueryOperations createForQuery(String thingType, String templateId, St
7070
TableSafeMetricBuilder.of(context.getMetricBuilder()),
7171
context.getSettings(),
7272
registry,
73-
operations.database());
73+
operations.database(),
74+
operations.schema());
7475
}
7576

7677
@Override

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/thing/TimescaleDBRowModeQueryOperations.java

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -28,8 +28,6 @@
2828
import org.hswebframework.web.api.crud.entity.PagerResult;
2929
import org.hswebframework.web.api.crud.entity.QueryParamEntity;
3030
import org.hswebframework.web.crud.query.QueryHelper;
31-
import org.jetlinks.core.metadata.EventMetadata;
32-
import org.jetlinks.core.things.ThingsRegistry;
3331
import org.jetlinks.community.things.data.AggregationRequest;
3432
import org.jetlinks.community.things.data.PropertyAggregation;
3533
import org.jetlinks.community.things.data.ThingsDataConstants;
@@ -41,6 +39,7 @@
4139
import org.jetlinks.community.timeseries.TimeSeriesData;
4240
import org.jetlinks.community.timeseries.query.Aggregation;
4341
import org.jetlinks.community.timeseries.query.AggregationData;
42+
import org.jetlinks.core.things.ThingsRegistry;
4443
import org.jetlinks.reactor.ql.utils.CastUtils;
4544
import org.slf4j.Logger;
4645
import reactor.core.publisher.Flux;
@@ -49,21 +48,21 @@
4948
import java.util.*;
5049
import java.util.function.Function;
5150

52-
import static org.jetlinks.community.timescaledb.thing.TimescaleDBColumnModeQueryOperations.doAggregation0;
53-
5451
@Slf4j
5552
public class TimescaleDBRowModeQueryOperations extends RowModeQueryOperationsBase {
5653
private final DatabaseOperator database;
57-
54+
private final String schema;
5855
public TimescaleDBRowModeQueryOperations(String thingType,
5956
String thingTemplateId,
6057
String thingId,
6158
MetricBuilder metricBuilder,
6259
DataSettings settings,
6360
ThingsRegistry registry,
64-
DatabaseOperator database) {
61+
DatabaseOperator database,
62+
String schema) {
6563
super(thingType, thingTemplateId, thingId, metricBuilder, settings, registry);
6664
this.database = database;
65+
this.schema=schema;
6766
}
6867

6968
@Override
@@ -117,7 +116,8 @@ protected Flux<AggregationData> doAggregation(String metric,
117116
if (request.getInterval() != null) {
118117
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(
119118
request.getFrom().getTime(),
120-
request.getInterval()
119+
request.getInterval(),
120+
schema
121121
);
122122

123123
query.groupBy(column);

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/thing/TimescaleDBRowModeStrategy.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,13 @@
1515
*/
1616
package org.jetlinks.community.timescaledb.thing;
1717

18-
import org.jetlinks.core.things.ThingsRegistry;
1918
import org.jetlinks.community.things.data.AbstractThingDataRepositoryStrategy;
2019
import org.jetlinks.community.things.data.operations.DDLOperations;
2120
import org.jetlinks.community.things.data.operations.QueryOperations;
2221
import org.jetlinks.community.things.data.operations.SaveOperations;
2322
import org.jetlinks.community.things.data.operations.TableSafeMetricBuilder;
2423
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
24+
import org.jetlinks.core.things.ThingsRegistry;
2525

2626
public class TimescaleDBRowModeStrategy extends AbstractThingDataRepositoryStrategy {
2727

@@ -65,7 +65,8 @@ protected QueryOperations createForQuery(String thingType, String templateId, St
6565
TableSafeMetricBuilder.of(context.getMetricBuilder()),
6666
context.getSettings(),
6767
registry,
68-
operations.database());
68+
operations.database(),
69+
operations.schema());
6970
}
7071

7172
@Override

jetlinks-components/timescaledb-component/src/main/java/org/jetlinks/community/timescaledb/timeseries/TimescaleDBTimeSeriesService.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,14 +24,14 @@
2424
import org.hswebframework.ezorm.rdb.operator.dml.query.SelectColumn;
2525
import org.hswebframework.web.bean.FastBeanCopier;
2626
import org.hswebframework.web.id.IDGenerator;
27-
import org.jetlinks.core.utils.Reactors;
2827
import org.jetlinks.community.things.data.ThingsDataConstants;
2928
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
3029
import org.jetlinks.community.timescaledb.TimescaleDBUtils;
3130
import org.jetlinks.community.timeseries.TimeSeriesData;
3231
import org.jetlinks.community.timeseries.TimeSeriesService;
3332
import org.jetlinks.community.timeseries.query.*;
3433
import org.jetlinks.community.timeseries.utils.TimeSeriesUtils;
34+
import org.jetlinks.core.utils.Reactors;
3535
import org.jetlinks.reactor.ql.utils.CastUtils;
3636
import org.reactivestreams.Publisher;
3737
import org.slf4j.Logger;
@@ -113,7 +113,7 @@ public Flux<AggregationData> aggregation(AggregationQueryParam param) {
113113
for (Group group : groups) {
114114
if (group instanceof TimeGroup) {
115115
_timeGroup = ((TimeGroup) group);
116-
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(startWith, _timeGroup.getInterval());
116+
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(startWith, _timeGroup.getInterval(), operations.schema());
117117
column.setColumn(ThingsDataConstants.COLUMN_TIMESTAMP);
118118
column.setAlias(group.getAlias());
119119
query.select(column);

0 commit comments

Comments
 (0)