@ -150,7 +150,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
if ( strategy . sendWsUpdate ( ) ) {
if ( strategy . sendWsUpdate ( ) ) {
addWsCallback ( saveFuture , success - > onTimeSeriesUpdate ( tenantId , entityId , request . getEntries ( ) ) ) ;
addWsCallback ( saveFuture , success - > onTimeSeriesUpdate ( tenantId , entityId , request . getEntries ( ) ) ) ;
}
}
if ( strategy . saveLatest ( ) ) {
if ( strategy . saveLatest ( ) & & entityId . getEntityType ( ) . isOneOf ( EntityType . DEVICE , EntityType . ASSET ) ) {
addMainCallback ( saveFuture , __ - > copyLatestToEntityViews ( tenantId , entityId , request . getEntries ( ) ) ) ;
addMainCallback ( saveFuture , __ - > copyLatestToEntityViews ( tenantId , entityId , request . getEntries ( ) ) ) ;
}
}
return saveFuture ;
return saveFuture ;
@ -207,58 +207,56 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}
}
private void copyLatestToEntityViews ( TenantId tenantId , EntityId entityId , List < TsKvEntry > ts ) {
private void copyLatestToEntityViews ( TenantId tenantId , EntityId entityId , List < TsKvEntry > ts ) {
if ( entityId . getEntityType ( ) . isOneOf ( EntityType . DEVICE , EntityType . ASSET ) ) {
Futures . addCallback ( tbEntityViewService . findEntityViewsByTenantIdAndEntityIdAsync ( tenantId , entityId ) ,
Futures . addCallback ( tbEntityViewService . findEntityViewsByTenantIdAndEntityIdAsync ( tenantId , entityId ) ,
new FutureCallback < > ( ) {
new FutureCallback < > ( ) {
@Override
@Override
public void onSuccess ( @Nullable List < EntityView > result ) {
public void onSuccess ( @Nullable List < EntityView > result ) {
if ( result ! = null & & ! result . isEmpty ( ) ) {
if ( result ! = null & & ! result . isEmpty ( ) ) {
Map < String , List < TsKvEntry > > tsMap = new HashMap < > ( ) ;
Map < String , List < TsKvEntry > > tsMap = new HashMap < > ( ) ;
for ( TsKvEntry entry : ts ) {
for ( TsKvEntry entry : ts ) {
tsMap . computeIfAbsent ( entry . getKey ( ) , s - > new ArrayList < > ( ) ) . add ( entry ) ;
tsMap . computeIfAbsent ( entry . getKey ( ) , s - > new ArrayList < > ( ) ) . add ( entry ) ;
}
}
for ( EntityView entityView : result ) {
for ( EntityView entityView : result ) {
List < String > keys = entityView . getKeys ( ) ! = null & & entityView . getKeys ( ) . getTimeseries ( ) ! = null ?
List < String > keys = entityView . getKeys ( ) ! = null & & entityView . getKeys ( ) . getTimeseries ( ) ! = null ?
entityView . getKeys ( ) . getTimeseries ( ) : new ArrayList < > ( tsMap . keySet ( ) ) ;
entityView . getKeys ( ) . getTimeseries ( ) : new ArrayList < > ( tsMap . keySet ( ) ) ;
List < TsKvEntry > entityViewLatest = new ArrayList < > ( ) ;
List < TsKvEntry > entityViewLatest = new ArrayList < > ( ) ;
long startTs = entityView . getStartTimeMs ( ) ;
long startTs = entityView . getStartTimeMs ( ) ;
long endTs = entityView . getEndTimeMs ( ) = = 0 ? Long . MAX_VALUE : entityView . getEndTimeMs ( ) ;
long endTs = entityView . getEndTimeMs ( ) = = 0 ? Long . MAX_VALUE : entityView . getEndTimeMs ( ) ;
for ( String key : keys ) {
for ( String key : keys ) {
List < TsKvEntry > entries = tsMap . get ( key ) ;
List < TsKvEntry > entries = tsMap . get ( key ) ;
if ( entries ! = null ) {
if ( entries ! = null ) {
Optional < TsKvEntry > tsKvEntry = entries . stream ( )
Optional < TsKvEntry > tsKvEntry = entries . stream ( )
. filter ( entry - > entry . getTs ( ) > startTs & & entry . getTs ( ) < = endTs )
. filter ( entry - > entry . getTs ( ) > startTs & & entry . getTs ( ) < = endTs )
. max ( Comparator . comparingLong ( TsKvEntry : : getTs ) ) ;
. max ( Comparator . comparingLong ( TsKvEntry : : getTs ) ) ;
tsKvEntry . ifPresent ( entityViewLatest : : add ) ;
tsKvEntry . ifPresent ( entityViewLatest : : add ) ;
}
}
if ( ! entityViewLatest . isEmpty ( ) ) {
saveTimeseries ( TimeseriesSaveRequest . builder ( )
. tenantId ( tenantId )
. entityId ( entityView . getId ( ) )
. entries ( entityViewLatest )
. strategy ( TimeseriesSaveRequest . Strategy . LATEST_AND_WS )
. callback ( new FutureCallback < > ( ) {
@Override
public void onSuccess ( @Nullable Void tmp ) { }
@Override
public void onFailure ( Throwable t ) {
log . error ( "[{}][{}] Failed to save entity view latest timeseries: {}" , tenantId , entityView . getId ( ) , entityViewLatest , t ) ;
}
} )
. build ( ) ) ;
}
}
}
}
if ( ! entityViewLatest . isEmpty ( ) ) {
saveTimeseries ( TimeseriesSaveRequest . builder ( )
. tenantId ( tenantId )
. entityId ( entityView . getId ( ) )
. entries ( entityViewLatest )
. strategy ( TimeseriesSaveRequest . Strategy . LATEST_AND_WS )
. callback ( new FutureCallback < > ( ) {
@Override
public void onSuccess ( @Nullable Void tmp ) { }
@Override
public void onFailure ( Throwable t ) {
log . error ( "[{}][{}] Failed to save entity view latest timeseries: {}" , tenantId , entityView . getId ( ) , entityViewLatest , t ) ;
}
} )
. build ( ) ) ;
}
}
}
}
}
}
@Override
@Override
public void onFailure ( Throwable t ) {
public void onFailure ( Throwable t ) {
log . error ( "Error while finding entity views by tenantId and entityId" , t ) ;
log . error ( "Error while finding entity views by tenantId and entityId" , t ) ;
}
}
} , MoreExecutors . directExecutor ( ) ) ;
} , MoreExecutors . directExecutor ( ) ) ;
}
}
}
private void onAttributesUpdate ( TenantId tenantId , EntityId entityId , String scope , List < AttributeKvEntry > attributes , boolean notifyDevice ) {
private void onAttributesUpdate ( TenantId tenantId , EntityId entityId , String scope , List < AttributeKvEntry > attributes , boolean notifyDevice ) {
@ -313,7 +311,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}
}
private < S > void addMainCallback ( ListenableFuture < S > saveFuture , Consumer < S > onSuccess ) {
private < S > void addMainCallback ( ListenableFuture < S > saveFuture , Consumer < S > onSuccess ) {
DonAsynchron . with Callback( saveFuture , onSuccess , null , tsCallBackExecutor ) ;
addMain Callback( saveFuture , onSuccess , null ) ;
}
}
private < S > void addMainCallback ( ListenableFuture < S > saveFuture , Consumer < S > onSuccess , Consumer < Throwable > onFailure ) {
private < S > void addMainCallback ( ListenableFuture < S > saveFuture , Consumer < S > onSuccess , Consumer < Throwable > onFailure ) {