|
|
|
@ -3,6 +3,7 @@ package org.gcube.accounting.aggregator.persistence;
|
|
|
|
|
import java.io.Serializable;
|
|
|
|
|
import java.sql.Connection;
|
|
|
|
|
import java.sql.DriverManager;
|
|
|
|
|
import java.sql.ResultSet;
|
|
|
|
|
import java.sql.SQLException;
|
|
|
|
|
import java.sql.Statement;
|
|
|
|
|
import java.text.SimpleDateFormat;
|
|
|
|
@ -11,11 +12,15 @@ import java.util.Calendar;
|
|
|
|
|
import java.util.Date;
|
|
|
|
|
import java.util.List;
|
|
|
|
|
import java.util.TimeZone;
|
|
|
|
|
import java.util.UUID;
|
|
|
|
|
|
|
|
|
|
import org.gcube.accounting.aggregator.aggregation.AggregationInfo;
|
|
|
|
|
import org.gcube.accounting.aggregator.aggregation.AggregationType;
|
|
|
|
|
import org.gcube.accounting.aggregator.status.AggregationState;
|
|
|
|
|
import org.gcube.accounting.aggregator.status.AggregationStateEvent;
|
|
|
|
|
import org.gcube.accounting.aggregator.status.AggregationStatus;
|
|
|
|
|
import org.gcube.accounting.aggregator.utility.Constant;
|
|
|
|
|
import org.gcube.accounting.aggregator.utility.Utility;
|
|
|
|
|
import org.gcube.accounting.analytics.persistence.AccountingPersistenceBackendQueryConfiguration;
|
|
|
|
|
import org.gcube.accounting.analytics.persistence.postgresql.AccountingPersistenceQueryPostgreSQL;
|
|
|
|
|
import org.gcube.accounting.persistence.AccountingPersistenceConfiguration;
|
|
|
|
@ -29,19 +34,32 @@ public class PostgreSQLConnector extends AccountingPersistenceQueryPostgreSQL {
|
|
|
|
|
private static final String UTC_TIME_ZONE = "UTC";
|
|
|
|
|
public static final TimeZone DEFAULT_TIME_ZONE = TimeZone.getTimeZone(UTC_TIME_ZONE);
|
|
|
|
|
|
|
|
|
|
public PostgreSQLConnector() throws Exception {
|
|
|
|
|
protected Connection connection;
|
|
|
|
|
|
|
|
|
|
private static PostgreSQLConnector postgreSQLConnector;
|
|
|
|
|
|
|
|
|
|
public static PostgreSQLConnector getPostgreSQLConnector() throws Exception {
|
|
|
|
|
if(postgreSQLConnector == null) {
|
|
|
|
|
postgreSQLConnector = new PostgreSQLConnector();
|
|
|
|
|
}
|
|
|
|
|
return postgreSQLConnector;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private PostgreSQLConnector() throws Exception {
|
|
|
|
|
this.configuration = new AccountingPersistenceBackendQueryConfiguration(AccountingPersistenceQueryPostgreSQL.class);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
protected Connection getConnection() throws Exception {
|
|
|
|
|
Class.forName("org.postgresql.Driver");
|
|
|
|
|
String url = configuration.getProperty(AccountingPersistenceQueryPostgreSQL.URL_PROPERTY_KEY);
|
|
|
|
|
String username = configuration.getProperty(AccountingPersistenceConfiguration.USERNAME_PROPERTY_KEY);
|
|
|
|
|
String password = configuration.getProperty(AccountingPersistenceConfiguration.PASSWORD_PROPERTY_KEY);
|
|
|
|
|
if(connection==null) {
|
|
|
|
|
Class.forName("org.postgresql.Driver");
|
|
|
|
|
String url = configuration.getProperty(AccountingPersistenceQueryPostgreSQL.URL_PROPERTY_KEY);
|
|
|
|
|
String username = configuration.getProperty(AccountingPersistenceConfiguration.USERNAME_PROPERTY_KEY);
|
|
|
|
|
String password = configuration.getProperty(AccountingPersistenceConfiguration.PASSWORD_PROPERTY_KEY);
|
|
|
|
|
|
|
|
|
|
Connection connection = DriverManager.getConnection(url, username, password);
|
|
|
|
|
logger.trace("Database {} opened successfully", url);
|
|
|
|
|
connection.setAutoCommit(false);
|
|
|
|
|
connection = DriverManager.getConnection(url, username, password);
|
|
|
|
|
logger.trace("Database {} opened successfully", url);
|
|
|
|
|
connection.setAutoCommit(false);
|
|
|
|
|
}
|
|
|
|
|
return connection;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -168,31 +186,113 @@ public class PostgreSQLConnector extends AccountingPersistenceQueryPostgreSQL {
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public void upsertAggregationStatus(AggregationStatus aggregationStatus) throws Exception {
|
|
|
|
|
Connection connection = getConnection();
|
|
|
|
|
Statement statement = connection.createStatement();
|
|
|
|
|
Statement statement = getConnection().createStatement();
|
|
|
|
|
String sqlCommand = getInsertAggregationStatusQuery(aggregationStatus, true);
|
|
|
|
|
statement.executeUpdate(sqlCommand);
|
|
|
|
|
sqlCommand = getInsertAggregationStateQuery(aggregationStatus);
|
|
|
|
|
statement.executeUpdate(sqlCommand);
|
|
|
|
|
statement.close();
|
|
|
|
|
connection.commit();
|
|
|
|
|
connection.close();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// private Calendar getCalendar(ResultSet resultSet, String columnLabel) throws SQLException {
|
|
|
|
|
// Date date = resultSet.getDate(columnLabel);
|
|
|
|
|
// Calendar calendar = Calendar.getInstance();
|
|
|
|
|
// calendar.setTime(date);
|
|
|
|
|
// return calendar;
|
|
|
|
|
// }
|
|
|
|
|
|
|
|
|
|
private AggregationStatus getAggregationStatusFromResultSet(ResultSet resultSet) throws Exception {
|
|
|
|
|
String recordType = resultSet.getString("record_type");
|
|
|
|
|
String aggregationTypeString = resultSet.getString("aggregation_type");
|
|
|
|
|
AggregationType aggregationType = AggregationType.valueOf(aggregationTypeString);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Date aggregationStartDate = resultSet.getDate("aggregation_start_date");
|
|
|
|
|
Date aggregationEndDate = resultSet.getDate("aggregation_end_date");
|
|
|
|
|
AggregationInfo aggregationInfo = new AggregationInfo(recordType, aggregationType, aggregationStartDate, aggregationEndDate);
|
|
|
|
|
|
|
|
|
|
AggregationStatus aggregationStatus = new AggregationStatus(aggregationInfo);
|
|
|
|
|
|
|
|
|
|
UUID uuid = UUID.fromString(resultSet.getString("id"));
|
|
|
|
|
aggregationStatus.setUUID(uuid);
|
|
|
|
|
|
|
|
|
|
int originalRecordsNumber = resultSet.getInt("original_records_number");
|
|
|
|
|
int aggregatedRecordsNumber = resultSet.getInt("aggregated_records_number");
|
|
|
|
|
int malformedRecordNumber = resultSet.getInt("malformed_records_number");
|
|
|
|
|
aggregationStatus.setRecordNumbers(originalRecordsNumber, aggregatedRecordsNumber, malformedRecordNumber);
|
|
|
|
|
|
|
|
|
|
String context = resultSet.getString("context");
|
|
|
|
|
aggregationStatus.setContext(context);
|
|
|
|
|
|
|
|
|
|
String current_aggregation_state = resultSet.getString("current_aggregation_state");
|
|
|
|
|
AggregationState aggregationState = AggregationState.valueOf(current_aggregation_state);
|
|
|
|
|
aggregationStatus.setAggregationState(aggregationState);
|
|
|
|
|
|
|
|
|
|
Date last_update_time = resultSet.getDate("last_update_time");
|
|
|
|
|
Calendar lastUpdateTime = Calendar.getInstance();
|
|
|
|
|
lastUpdateTime.setTime(last_update_time);
|
|
|
|
|
aggregationStatus.setLastUpdateTime(lastUpdateTime);
|
|
|
|
|
|
|
|
|
|
return aggregationStatus;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public AggregationStatus getLast(String recordType, AggregationType aggregationType, Date aggregationStartDate, Date aggregationEndDate) throws Exception{
|
|
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
* SELECT *
|
|
|
|
|
* FROM aggregation_status
|
|
|
|
|
* SELECT * FROM aggregation_status
|
|
|
|
|
*
|
|
|
|
|
* WHERE
|
|
|
|
|
* `aggregationInfo`.`recordType` = "ServiceUsageRecord" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationType` = "DAILY" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationStartDate` >= "2017-05-01 00:00:00.000 +0000"
|
|
|
|
|
* `aggregationInfo`.`aggregationStartDate` <= "2017-05-31 00:00:00.000 +0000"
|
|
|
|
|
* ORDER BY `aggregationInfo`.`aggregationStartDate` DESC LIMIT 1
|
|
|
|
|
* record_type = 'ServiceUsageRecord' AND
|
|
|
|
|
* aggregation_type = 'DAILY' AND
|
|
|
|
|
* aggregation_start_date >= '2017-05-01 00:00:00.000 +0000' AND
|
|
|
|
|
* aggregation_start_date <= '2017-05-31 00:00:00.000 +0000'
|
|
|
|
|
*
|
|
|
|
|
* ORDER BY last_update_time DESC LIMIT 1
|
|
|
|
|
*
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
return null;
|
|
|
|
|
StringBuffer stringBuffer = new StringBuffer();
|
|
|
|
|
stringBuffer.append("SELECT * ");
|
|
|
|
|
stringBuffer.append("FROM aggregation_status");
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" WHERE ");
|
|
|
|
|
stringBuffer.append("record_type = ");
|
|
|
|
|
stringBuffer.append(getValue(recordType));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_type = ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationType.name()));
|
|
|
|
|
|
|
|
|
|
if(aggregationStartDate!=null && aggregationEndDate!=null) {
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date >= ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationStartDate));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date <= ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationEndDate));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" ORDER BY ");
|
|
|
|
|
stringBuffer.append("last_update_time DESC LIMIT 1");
|
|
|
|
|
|
|
|
|
|
Connection connection = getConnection();
|
|
|
|
|
Statement statement = connection.createStatement();
|
|
|
|
|
|
|
|
|
|
String sqlCommand = stringBuffer.toString();
|
|
|
|
|
|
|
|
|
|
logger.trace("Going to request the following query: {}", sqlCommand);
|
|
|
|
|
ResultSet resultSet = statement.executeQuery(sqlCommand);
|
|
|
|
|
|
|
|
|
|
AggregationStatus aggregationStatus = null;
|
|
|
|
|
|
|
|
|
|
while(resultSet.next()) {
|
|
|
|
|
aggregationStatus = getAggregationStatusFromResultSet(resultSet);
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return aggregationStatus;
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@ -204,20 +304,65 @@ public class PostgreSQLConnector extends AccountingPersistenceQueryPostgreSQL {
|
|
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
* SELECT *
|
|
|
|
|
* FROM AccountingManager
|
|
|
|
|
* FROM aggregation_status
|
|
|
|
|
* WHERE
|
|
|
|
|
* `aggregationState` != "COMPLETED" AND
|
|
|
|
|
* `lastUpdateTime` < "2017-07-31 09:31:10.984 +0000" AND
|
|
|
|
|
* `aggregationInfo`.`recordType` = "ServiceUsageRecord" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationType` = "DAILY" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationStartDate` >= "2017-05-01 00:00:00.000 +0000"
|
|
|
|
|
* `aggregationInfo`.`aggregationStartDate` <= "2017-05-31 00:00:00.000 +0000"
|
|
|
|
|
* aggregation_state != "COMPLETED" AND
|
|
|
|
|
* last_update_time < "2017-07-31 09:31:10.984 +0000" AND
|
|
|
|
|
* record_type = "ServiceUsageRecord" AND
|
|
|
|
|
* aggregation_type` = "DAILY" AND
|
|
|
|
|
* aggregation_start_date >= "2017-05-01 00:00:00.000 +0000"
|
|
|
|
|
* aggregation_start_date <= "2017-05-31 00:00:00.000 +0000"
|
|
|
|
|
*
|
|
|
|
|
* ORDER BY `aggregationInfo`.`aggregationStartDate` ASC
|
|
|
|
|
* ORDER BY aggregation_start_date ASC
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
StringBuffer stringBuffer = new StringBuffer();
|
|
|
|
|
stringBuffer.append("SELECT * ");
|
|
|
|
|
stringBuffer.append("FROM aggregation_status");
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" WHERE ");
|
|
|
|
|
stringBuffer.append("current_aggregation_state != ");
|
|
|
|
|
stringBuffer.append(getValue(AggregationState.COMPLETED));
|
|
|
|
|
|
|
|
|
|
Calendar now = Utility.getUTCCalendarInstance();
|
|
|
|
|
now.add(Constant.CALENDAR_FIELD_TO_SUBSTRACT_TO_CONSIDER_UNTERMINATED, -Constant.UNIT_TO_SUBSTRACT_TO_CONSIDER_UNTERMINATED);
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("last_update_time < ");
|
|
|
|
|
stringBuffer.append(getValue(now));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("record_type = ");
|
|
|
|
|
stringBuffer.append(getValue(recordType));
|
|
|
|
|
|
|
|
|
|
if(aggregationStartDate!=null && aggregationEndDate!=null) {
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date >= ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationStartDate));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date <= ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationEndDate));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" ORDER BY ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date ASC");
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Connection connection = getConnection();
|
|
|
|
|
Statement statement = connection.createStatement();
|
|
|
|
|
|
|
|
|
|
String sqlCommand = stringBuffer.toString();
|
|
|
|
|
|
|
|
|
|
logger.trace("Going to request the following query: {}", sqlCommand);
|
|
|
|
|
ResultSet resultSet = statement.executeQuery(sqlCommand);
|
|
|
|
|
|
|
|
|
|
List<AggregationStatus> aggregationStatuses = new ArrayList<>();
|
|
|
|
|
|
|
|
|
|
while(resultSet.next()) {
|
|
|
|
|
AggregationStatus aggregationStatus = getAggregationStatusFromResultSet(resultSet);
|
|
|
|
|
aggregationStatuses.add(aggregationStatus);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return aggregationStatuses;
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
@ -226,15 +371,33 @@ public class PostgreSQLConnector extends AccountingPersistenceQueryPostgreSQL {
|
|
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
* SELECT *
|
|
|
|
|
* FROM AccountingManager
|
|
|
|
|
* ORDER BY `aggregationInfo`.`aggregationStartDate` ASC
|
|
|
|
|
* FROM aggregation_status
|
|
|
|
|
* ORDER BY aggregation_start_date ASC
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
StringBuffer stringBuffer = new StringBuffer();
|
|
|
|
|
stringBuffer.append("SELECT * ");
|
|
|
|
|
stringBuffer.append("FROM aggregation_status");
|
|
|
|
|
stringBuffer.append("ORDER BY aggregation_start_date ASC");
|
|
|
|
|
|
|
|
|
|
Connection connection = getConnection();
|
|
|
|
|
Statement statement = connection.createStatement();
|
|
|
|
|
|
|
|
|
|
String sqlCommand = stringBuffer.toString();
|
|
|
|
|
|
|
|
|
|
logger.trace("Going to request the following query: {}", sqlCommand);
|
|
|
|
|
ResultSet resultSet = statement.executeQuery(sqlCommand);
|
|
|
|
|
|
|
|
|
|
List<AggregationStatus> aggregationStatuses = new ArrayList<>();
|
|
|
|
|
|
|
|
|
|
while(resultSet.next()) {
|
|
|
|
|
AggregationStatus aggregationStatus = getAggregationStatusFromResultSet(resultSet);
|
|
|
|
|
aggregationStatuses.add(aggregationStatus);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return aggregationStatuses;
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@ -242,14 +405,50 @@ public class PostgreSQLConnector extends AccountingPersistenceQueryPostgreSQL {
|
|
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
* SELECT *
|
|
|
|
|
* FROM AccountingManager
|
|
|
|
|
* FROM aggregation_status
|
|
|
|
|
* WHERE
|
|
|
|
|
* `aggregationInfo`.`recordType` = "ServiceUsageRecord" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationType` = "DAILY" AND
|
|
|
|
|
* `aggregationInfo`.`aggregationStartDate` = "2017-06-24 00:00:00.000 +0000"
|
|
|
|
|
* record_type = "ServiceUsageRecord" AND
|
|
|
|
|
* aggregation_type = "DAILY" AND
|
|
|
|
|
* aggregation_start_date` = "2017-06-24 00:00:00.000 +0000"
|
|
|
|
|
* ORDER BY aggregation_start_date DESC LIMIT 1
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
return null;
|
|
|
|
|
|
|
|
|
|
StringBuffer stringBuffer = new StringBuffer();
|
|
|
|
|
stringBuffer.append("SELECT * ");
|
|
|
|
|
stringBuffer.append("FROM aggregation_status");
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" WHERE ");
|
|
|
|
|
stringBuffer.append("record_type = ");
|
|
|
|
|
stringBuffer.append(getValue(recordType));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_type = ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationType.name()));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" AND ");
|
|
|
|
|
stringBuffer.append("aggregation_start_date = ");
|
|
|
|
|
stringBuffer.append(getValue(aggregationStartDate));
|
|
|
|
|
|
|
|
|
|
stringBuffer.append(" ORDER BY ");
|
|
|
|
|
stringBuffer.append("last_update_time DESC LIMIT 1");
|
|
|
|
|
|
|
|
|
|
Connection connection = getConnection();
|
|
|
|
|
Statement statement = connection.createStatement();
|
|
|
|
|
|
|
|
|
|
String sqlCommand = stringBuffer.toString();
|
|
|
|
|
|
|
|
|
|
logger.trace("Going to request the following query: {}", sqlCommand);
|
|
|
|
|
ResultSet resultSet = statement.executeQuery(sqlCommand);
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
AggregationStatus aggregationStatus = null;
|
|
|
|
|
while(resultSet.next()) {
|
|
|
|
|
aggregationStatus = getAggregationStatusFromResultSet(resultSet);
|
|
|
|
|
break;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return aggregationStatus;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|