Skip to content
Closed
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
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@ public interface TimescaleDBOperations {

DatabaseOperator database();

String functionSchema();

TimescaleDBDataWriter writer();

}
Original file line number Diff line number Diff line change
Expand Up @@ -36,8 +36,8 @@ public class TimescaleDBProperties {
//当sharedSpring未false时,使用此连接配置.
private R2dbcProperties r2dbc = new R2dbcProperties();

//数据库的schema
private String schema = "public";
//数据函数的schema
private String functionSchema = "public";

/**
* 写入缓冲区配置
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,16 @@
import org.hswebframework.ezorm.rdb.operator.dml.query.NativeSelectColumn;
import org.hswebframework.ezorm.rdb.operator.dml.query.SelectColumn;
import org.hswebframework.ezorm.rdb.supports.postgres.JsonbType;
import org.jetlinks.core.metadata.DataType;
import org.jetlinks.core.metadata.PropertyMetadata;
import org.jetlinks.core.metadata.types.ArrayType;
import org.jetlinks.core.metadata.types.ObjectType;
import org.jetlinks.community.Interval;
import org.jetlinks.community.things.data.ThingsDataConstants;
import org.jetlinks.community.things.utils.ThingsDatabaseUtils;
import org.jetlinks.community.timescaledb.metadata.JsonbValueCodec;
import org.jetlinks.community.timeseries.TimeSeriesData;
import org.jetlinks.community.timeseries.query.Aggregation;
import org.jetlinks.community.utils.ObjectMappers;
import org.jetlinks.core.metadata.DataType;
import org.jetlinks.core.metadata.PropertyMetadata;
import org.jetlinks.core.metadata.types.ArrayType;
import org.jetlinks.core.metadata.types.ObjectType;
import org.jetlinks.reactor.ql.utils.CastUtils;

public class TimescaleDBUtils {
Expand All @@ -40,15 +39,15 @@ public static String getTableName(String name) {
return ThingsDatabaseUtils.createTableName(name);
}

public static NativeSelectColumn createTimeGroupColumn(long startWith, Interval interval) {
public static NativeSelectColumn createTimeGroupColumn(long startWith, Interval interval, String schema) {

String unit = interval.getNumber().intValue() + " " + interval
.getUnit()
.name()
.toLowerCase();

return NativeSelectColumn
.of("time_bucket('" + unit + "',timestamp)");
.of(schema + "." + "time_bucket('" + unit + "',timestamp)");
}

public static TimeSeriesData convertToTimeSeriesData(Record record) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ public void init() {
database.addFeature(sqlExecutor);
database.addFeature(ReactiveSyncSqlExecutor.of(sqlExecutor));

RDBSchemaMetadata schema = TimescaleDBDialectProvider.GLOBAL.createSchema(properties.getSchema());
RDBSchemaMetadata schema = TimescaleDBDialectProvider.GLOBAL.createSchema(properties.getFunctionSchema());
database.addSchema(schema);
database.setCurrentSchema(schema);
this.database = DefaultDatabaseOperator.of(database);
Expand All @@ -75,7 +75,7 @@ public void init() {
}
RDBDataSourceProperties datasource = new RDBDataSourceProperties();
datasource.setType(RDBDataSourceProperties.Type.r2dbc);
datasource.setSchema(properties.getSchema());
datasource.setSchema(properties.getFunctionSchema());
datasource.setUsername(properties.getR2dbc().getUsername());
datasource.setPassword(properties.getR2dbc().getPassword());
datasource.setUrl(properties.getR2dbc().getUrl());
Expand All @@ -102,6 +102,11 @@ public DatabaseOperator database() {
return database;
}

@Override
public String functionSchema() {
return properties.getFunctionSchema();
}

@Override
public TimescaleDBDataWriter writer() {
return writer;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,26 +28,27 @@
import org.hswebframework.web.api.crud.entity.PagerResult;
import org.hswebframework.web.api.crud.entity.QueryParamEntity;
import org.hswebframework.web.crud.query.QueryHelper;
import org.jetlinks.core.metadata.EventMetadata;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AggregationRequest;
import org.jetlinks.community.things.data.PropertyAggregation;
import org.jetlinks.community.things.data.ThingsDataConstants;
import org.jetlinks.community.things.data.ThingsDataUtils;
import org.jetlinks.community.things.data.operations.ColumnModeQueryOperationsBase;
import org.jetlinks.community.things.data.operations.DataSettings;
import org.jetlinks.community.things.data.operations.MetricBuilder;
import org.jetlinks.community.things.data.operations.RowModeQueryOperationsBase;
import org.jetlinks.community.timescaledb.TimescaleDBUtils;
import org.jetlinks.community.timeseries.TimeSeriesData;
import org.jetlinks.community.timeseries.query.Aggregation;
import org.jetlinks.community.timeseries.query.AggregationData;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.reactor.ql.utils.CastUtils;
import org.slf4j.Logger;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import java.util.*;
import java.util.ArrayList;
import java.util.Map;
import java.util.NavigableMap;
import java.util.Objects;
import java.util.function.Function;

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

private final String schema;

public TimescaleDBColumnModeQueryOperations(String thingType,
String thingTemplateId,
String thingId,
MetricBuilder metricBuilder,
DataSettings settings,
ThingsRegistry registry,
DatabaseOperator database) {
DatabaseOperator database,
String schema) {
super(thingType, thingTemplateId, thingId, metricBuilder, settings, registry);
this.database = database;
this.schema = schema;
}

@Override
Expand Down Expand Up @@ -104,12 +109,13 @@ record -> mapper.apply(convertToTimeSeriesData(record))
protected Flux<AggregationData> doAggregation(String metric,
AggregationRequest request,
AggregationContext context) {
return doAggregation0(database, metric, request, context);
return doAggregation0(database, metric, schema, request, context);
}


static Flux<AggregationData> doAggregation0(DatabaseOperator database,
String metric,
String schema,
AggregationRequest request,
AggregationContext context) {
metric = TimescaleDBUtils.getTableName(metric);
Expand All @@ -122,7 +128,8 @@ static Flux<AggregationData> doAggregation0(DatabaseOperator database,
if (request.getInterval() != null) {
NativeSelectColumn column = createTimeGroupColumn(
request.getFrom().getTime(),
request.getInterval()
request.getInterval(),
schema
);
query.groupBy(column);
query.select(column);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,13 @@
*/
package org.jetlinks.community.timescaledb.thing;

import org.jetlinks.community.things.data.TableSafeMetricBuilder;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AbstractThingDataRepositoryStrategy;
import org.jetlinks.community.things.data.TableSafeMetricBuilder;
import org.jetlinks.community.things.data.operations.DDLOperations;
import org.jetlinks.community.things.data.operations.QueryOperations;
import org.jetlinks.community.things.data.operations.SaveOperations;
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
import org.jetlinks.core.things.ThingsRegistry;

public class TimescaleDBColumnModeStrategy extends AbstractThingDataRepositoryStrategy {

Expand Down Expand Up @@ -70,7 +70,8 @@ protected QueryOperations createForQuery(String thingType, String templateId, St
TableSafeMetricBuilder.of(context.getMetricBuilder()),
context.getSettings(),
registry,
operations.database());
operations.database(),
operations.functionSchema());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,8 +28,6 @@
import org.hswebframework.web.api.crud.entity.PagerResult;
import org.hswebframework.web.api.crud.entity.QueryParamEntity;
import org.hswebframework.web.crud.query.QueryHelper;
import org.jetlinks.core.metadata.EventMetadata;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AggregationRequest;
import org.jetlinks.community.things.data.PropertyAggregation;
import org.jetlinks.community.things.data.ThingsDataConstants;
Expand All @@ -41,6 +39,7 @@
import org.jetlinks.community.timeseries.TimeSeriesData;
import org.jetlinks.community.timeseries.query.Aggregation;
import org.jetlinks.community.timeseries.query.AggregationData;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.reactor.ql.utils.CastUtils;
import org.slf4j.Logger;
import reactor.core.publisher.Flux;
Expand All @@ -49,21 +48,21 @@
import java.util.*;
import java.util.function.Function;

import static org.jetlinks.community.timescaledb.thing.TimescaleDBColumnModeQueryOperations.doAggregation0;

@Slf4j
public class TimescaleDBRowModeQueryOperations extends RowModeQueryOperationsBase {
private final DatabaseOperator database;

private final String schema;
public TimescaleDBRowModeQueryOperations(String thingType,
String thingTemplateId,
String thingId,
MetricBuilder metricBuilder,
DataSettings settings,
ThingsRegistry registry,
DatabaseOperator database) {
DatabaseOperator database,
String schema) {
super(thingType, thingTemplateId, thingId, metricBuilder, settings, registry);
this.database = database;
this.schema=schema;
}

@Override
Expand Down Expand Up @@ -117,7 +116,8 @@ protected Flux<AggregationData> doAggregation(String metric,
if (request.getInterval() != null) {
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(
request.getFrom().getTime(),
request.getInterval()
request.getInterval(),
schema
);

query.groupBy(column);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,13 @@
*/
package org.jetlinks.community.timescaledb.thing;

import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AbstractThingDataRepositoryStrategy;
import org.jetlinks.community.things.data.operations.DDLOperations;
import org.jetlinks.community.things.data.operations.QueryOperations;
import org.jetlinks.community.things.data.operations.SaveOperations;
import org.jetlinks.community.things.data.operations.TableSafeMetricBuilder;
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
import org.jetlinks.core.things.ThingsRegistry;

public class TimescaleDBRowModeStrategy extends AbstractThingDataRepositoryStrategy {

Expand Down Expand Up @@ -65,7 +65,8 @@ protected QueryOperations createForQuery(String thingType, String templateId, St
TableSafeMetricBuilder.of(context.getMetricBuilder()),
context.getSettings(),
registry,
operations.database());
operations.database(),
operations.functionSchema());
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,14 +24,14 @@
import org.hswebframework.ezorm.rdb.operator.dml.query.SelectColumn;
import org.hswebframework.web.bean.FastBeanCopier;
import org.hswebframework.web.id.IDGenerator;
import org.jetlinks.core.utils.Reactors;
import org.jetlinks.community.things.data.ThingsDataConstants;
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
import org.jetlinks.community.timescaledb.TimescaleDBUtils;
import org.jetlinks.community.timeseries.TimeSeriesData;
import org.jetlinks.community.timeseries.TimeSeriesService;
import org.jetlinks.community.timeseries.query.*;
import org.jetlinks.community.timeseries.utils.TimeSeriesUtils;
import org.jetlinks.core.utils.Reactors;
import org.jetlinks.reactor.ql.utils.CastUtils;
import org.reactivestreams.Publisher;
import org.slf4j.Logger;
Expand Down Expand Up @@ -113,7 +113,7 @@ public Flux<AggregationData> aggregation(AggregationQueryParam param) {
for (Group group : groups) {
if (group instanceof TimeGroup) {
_timeGroup = ((TimeGroup) group);
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(startWith, _timeGroup.getInterval());
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(startWith, _timeGroup.getInterval(), operations.functionSchema());
column.setColumn(ThingsDataConstants.COLUMN_TIMESTAMP);
column.setAlias(group.getAlias());
query.select(column);
Expand Down