Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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 @@ -39,6 +39,11 @@ public class TimescaleDBProperties {
//数据库的schema
private String schema = "public";

/**
* TimescaleDB超表函数所在的位置
*/
private String functionSchema = "public";

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

默认应该为null. 为null时使用schema, 否则只配置了schema的情况下, 可能有问题.


/**
* 写入缓冲区配置
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,15 +40,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 functionSchema) {

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

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

括起来更好? 有的schema 有- 可能会有问题.

}

public static TimeSeriesData convertToTimeSeriesData(Record record) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
package org.jetlinks.community.timescaledb.configuration;

import org.jetlinks.community.timescaledb.TimescaleDBOperations;
import org.jetlinks.community.timescaledb.TimescaleDBProperties;
import org.jetlinks.community.timescaledb.timeseries.TimescaleDBTimeSeriesManager;
import org.jetlinks.community.timescaledb.timeseries.TimescaleDBTimeSeriesProperties;
import org.springframework.boot.autoconfigure.AutoConfiguration;
Expand All @@ -35,8 +36,9 @@ public class TimescaleDBTimeSeriesConfiguration {
@Bean
@Primary
public TimescaleDBTimeSeriesManager timescaleDBTimeSeriesManager(TimescaleDBOperations operations,
TimescaleDBTimeSeriesProperties properties) {
return new TimescaleDBTimeSeriesManager(properties, operations);
TimescaleDBTimeSeriesProperties properties,
TimescaleDBProperties timescaleDBProperties) {
return new TimescaleDBTimeSeriesManager(properties, operations,timescaleDBProperties);
}


Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
package org.jetlinks.community.timescaledb.metadata;

import lombok.AllArgsConstructor;
import lombok.Getter;
import org.hswebframework.ezorm.core.FeatureId;
import org.hswebframework.ezorm.core.FeatureType;
import org.hswebframework.ezorm.core.meta.Feature;

@AllArgsConstructor(staticName = "of")
@Getter
public class FunctionSchema implements Feature, FeatureType {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

建立直接传 TimescaleDBProperties 便于后续拓展


public static final FeatureId<FunctionSchema> ID = FeatureId.of("FunctionSchema");

private final String functionSchema;

@Override
public String getId() {
return ID.getId();
}

@Override
public String getName() {
return "FunctionSchema";
}

@Override
public FeatureType getType() {
return this;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,16 +23,11 @@

public class TimescaleDBCreateTableSqlBuilder extends CommonCreateTableSqlBuilder {

private String schema;

public TimescaleDBCreateTableSqlBuilder(String schema) {
this.schema = schema;
}

@Override
public SqlRequest build(RDBTableMetadata table) {
DefaultBatchSqlRequest sqlRequest = (DefaultBatchSqlRequest) super.build(table);


table.getFeature(CreateHypertable.ID)
.ifPresent(createHypertable -> sqlRequest.addBatch(createCreateHypertableSQL(table, createHypertable)));

Expand All @@ -47,9 +42,11 @@ private SqlRequest createCreateRetentionPolicySQL(RDBTableMetadata table, Create

String interval = createHypertable.getInterval().getNumber().intValue() + " "
+ createHypertable.getInterval().getUnit().name().toLowerCase();
String functionSchema = table.getFeatureNow(FunctionSchema.ID)
.getFunctionSchema();

return SqlRequests.of(
"SELECT "+ schema +".add_retention_policy( ? , INTERVAL '" + interval + "')",
"SELECT " + functionSchema + ".add_retention_policy( ? , INTERVAL '" + interval + "')",
table.getFullName()
);
}
Expand All @@ -58,9 +55,10 @@ private SqlRequest createCreateHypertableSQL(RDBTableMetadata table, CreateHyper

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

String functionSchema = table.getFeatureNow(FunctionSchema.ID)
.getFunctionSchema();
return SqlRequests.of(
"SELECT "+ schema +".create_hypertable( ? , ? , chunk_time_interval => INTERVAL '" + interval + "')",
"SELECT " + functionSchema + ".create_hypertable( ? , ? , chunk_time_interval => INTERVAL '" + interval + "')",
table.getFullName(),
table.getColumnNow(createHypertable.getColumn()).getName()
);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ public String getBindSymbol() {
@Override
public RDBSchemaMetadata createSchema(String name) {
PostgresqlSchemaMetadata schema = new PostgresqlSchemaMetadata(name);
schema.addFeature(new TimescaleDBCreateTableSqlBuilder(name));
schema.addFeature(new TimescaleDBCreateTableSqlBuilder());
schema.addFeature(new TimescaleDBAlterTableSqlBuilder());
DefaultValueCodecFactory codecFactory = new DefaultValueCodecFactory();
codecFactory
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
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.community.timescaledb.metadata.FunctionSchema;
import org.jetlinks.core.metadata.EventMetadata;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AggregationRequest;
Expand Down Expand Up @@ -122,7 +123,8 @@ static Flux<AggregationData> doAggregation0(DatabaseOperator database,
if (request.getInterval() != null) {
NativeSelectColumn column = createTimeGroupColumn(
request.getFrom().getTime(),
request.getInterval()
request.getInterval(),
database.getMetadata().getTable(metric).get().getFeatureNow(FunctionSchema.ID).getFunctionSchema()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

某些特殊情况下没创建表时会产生错误.

);
query.groupBy(column);
query.select(column);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
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.community.timescaledb.metadata.FunctionSchema;
import org.jetlinks.core.metadata.EventMetadata;
import org.jetlinks.core.things.ThingsRegistry;
import org.jetlinks.community.things.data.AggregationRequest;
Expand Down Expand Up @@ -117,7 +118,8 @@ protected Flux<AggregationData> doAggregation(String metric,
if (request.getInterval() != null) {
NativeSelectColumn column = TimescaleDBUtils.createTimeGroupColumn(
request.getFrom().getTime(),
request.getInterval()
request.getInterval(),
database.getMetadata().getTable(metric).get().getFeatureNow(FunctionSchema.ID).getFunctionSchema()
);

query.groupBy(column);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
import org.hswebframework.ezorm.rdb.codec.DateTimeCodec;
import org.hswebframework.ezorm.rdb.metadata.RDBIndexMetadata;
import org.hswebframework.ezorm.rdb.operator.ddl.TableBuilder;
import org.jetlinks.community.timescaledb.TimescaleDBProperties;
import org.jetlinks.community.timescaledb.metadata.FunctionSchema;
import org.jetlinks.core.metadata.PropertyMetadata;
import org.jetlinks.community.Interval;
import org.jetlinks.community.things.data.ThingsDataConstants;
Expand Down Expand Up @@ -48,6 +50,8 @@ public class TimescaleDBTimeSeriesManager implements TimeSeriesManager {

private final TimescaleDBOperations operations;

private final TimescaleDBProperties timescaleDBProperties;

@Override
public TimeSeriesService getService(TimeSeriesMetric metric) {
return getService(metric.getId());
Expand Down Expand Up @@ -82,6 +86,7 @@ public Mono<Void> registerMetadata(TimeSeriesMetadata metadata) {
.ddl()
.createOrAlter(tableName)
.custom(table -> {
table.addFeature(FunctionSchema.of(timescaleDBProperties.getFunctionSchema()));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

建议创建database 时直接放到 database或者schema中

table.addFeature(new CreateHypertable(ThingsDataConstants.COLUMN_TIMESTAMP, properties.getChunkTimeInterval()));
Interval interval = properties.getRetentionPolicy(tableName);
if (interval != null && interval.getNumber().longValue() > 0) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.hswebframework.ezorm.rdb.operator.dml.query.SelectColumn;
import org.hswebframework.web.bean.FastBeanCopier;
import org.hswebframework.web.id.IDGenerator;
import org.jetlinks.community.timescaledb.metadata.FunctionSchema;
import org.jetlinks.core.utils.Reactors;
import org.jetlinks.community.things.data.ThingsDataConstants;
import org.jetlinks.community.timescaledb.TimescaleDBOperations;
Expand Down Expand Up @@ -113,7 +114,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.database().getMetadata().getTable(metric).get().getFeatureNow(FunctionSchema.ID).getFunctionSchema());
column.setColumn(ThingsDataConstants.COLUMN_TIMESTAMP);
column.setAlias(group.getAlias());
query.select(column);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ REDIS_DATABASE: 0 # redis 数据库索引

# timescalbedb相关配置
TIMESCALEDB_SCHEMA: public # timescaledb 数据库schema,默认public
TIMESCALEDB_FUNCTION_SCHEMA: public # 时序数据所用到相关函数所在的schema


# elasticsearch相关配置
Expand Down
1 change: 1 addition & 0 deletions jetlinks-standalone/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ timescaledb:
# username: postgres
# password: p@ssw0rd
schema: ${TIMESCALEDB_SCHEMA:${easyorm.default-schema}} # timescaledb的schema,默认public
function-schema: ${TIMESCALEDB_FUNCTION_SCHEMA:${timescaledb.schema}} #时序数据所用到相关函数所在的schema
time-series:
enabled: true
retention-policies:
Expand Down