Skip to content
Open
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 @@ -173,6 +173,9 @@ public class WSCalculationConfiguration {

@Value("${egov.watercalculatorservice.createdemand.topic}")
private String createDemand;

@Value("${egov.watercalculatorservice.createsingledemand.topic}")
private String createSingleDemand;

@Value("${ws.demand.based.batch.size}")
private Integer batchSize;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ public class DemandGenerationConsumer {
* Exception e){ log.error("KAFKA_PROCESS_ERROR", e); } }
*/
@KafkaListener(
topics = "${egov.watercalculatorservice.createdemand.topic}",
topics = {"${egov.watercalculatorservice.createdemand.topic}","${egov.watercalculatorservice.createsingledemand.topic}"},
containerFactory = "kafkaListenerContainerFactoryBatch",
concurrency = "${egov.watercalculatorservice.listener.concurrency}"
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ public List<String> getConnectionsByStatus(String string, String status) {

/**
* executes query to update bill scheduler status
* @param billIds
// * @param billIds
*/
public void updateBillSchedularStatus(String schedulerId, StatusEnum status) {

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,7 @@ public interface WSCalculationDao {

List<WaterConnection> getConnection(String tenantId, String consumerCode,String connectionType,Long fromDate, Long toDate);
List<String> getLocalityList(String tenantId, String locality);
List<WaterDetails> getConnectionsNoListforsingledemand(String tenantId, String connectionType, Long taxPeriodFrom,
List<WaterDetails> getConnectionsNoListforsingledemand(String tenantId, String locality, String connectionType, Long taxPeriodFrom,
Long taxPeriodTo, String cone);

List<String> fetchUsageCategory(String consumerCodes);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -379,14 +379,20 @@ public void updateBillStatus(List<String> consumerCodes, String businessService,
jdbcTemplate.update(queryStr, preparedStmtList.toArray());
}
@Override
public List<WaterDetails> getConnectionsNoListforsingledemand(String tenantId, String connectionType, Long taxPeriodFrom,
public List<WaterDetails> getConnectionsNoListforsingledemand(String tenantId, String locality, String connectionType, Long taxPeriodFrom,
Long taxPeriodTo, String cone) {

String query;
List<Object> preparedStatement = new ArrayList<>();
String query = queryBuilder.getConnectionNumberList(tenantId, connectionType,
WSCalculationConstant.ACTIVE_CONNECTION, taxPeriodFrom, taxPeriodTo, cone, preparedStatement);
log.info("preparedStatement: " + preparedStatement + " connection type: " + connectionType
+ " connection list : " + query);
if(locality!=null && !locality.trim().isEmpty()){
query = queryBuilder.getConnectionNumberList(tenantId, locality, connectionType,
WSCalculationConstant.ACTIVE_CONNECTION, taxPeriodFrom, taxPeriodTo, cone, preparedStatement);
}
else {
query = queryBuilder.getConnectionNumberList(tenantId, null, connectionType,
WSCalculationConstant.ACTIVE_CONNECTION, taxPeriodFrom, taxPeriodTo, cone, preparedStatement);
log.info("preparedStatement: " + preparedStatement + " connection type: " + connectionType
+ " connection list : " + query);
}
return jdbcTemplate.query(query, preparedStatement.toArray(), demandSchedulerRowMapper);

}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ public String getMeterId(String connectionNo, Long lastReadingDate, Long current
/**
* Bill expire query builder
*
* @param billIds
// * @param billIds
* @param preparedStmtList
*/
public String getBillSchedulerUpdateQuery(String schedulerId, List<Object> preparedStmtList) {
Expand Down Expand Up @@ -355,9 +355,10 @@ public String getLocalityListWithBatch(String tenantId, String batchCode, List<O
return query.toString();
}

public String getConnectionNumberList(String tenantId, String connectionType, String status, Long taxPeriodFrom,
public String getConnectionNumberList(String tenantId, String locality, String connectionType, String status, Long taxPeriodFrom,
Long taxPeriodTo, String cone, List<Object> preparedStatement) {
StringBuilder query = new StringBuilder(connectionNoListQuery);
query.append(" INNER JOIN eg_pt_property p ON conn.property_id = p.propertyid INNER JOIN eg_pt_address a ON a.propertyid = p.id ");

// Add connection type
addClauseIfRequired(preparedStatement, query);
Expand All @@ -380,6 +381,11 @@ public String getConnectionNumberList(String tenantId, String connectionType, St
query.append(" conn.tenantid = ? ");
preparedStatement.add(tenantId);

if (locality != null) {
addClauseIfRequired(preparedStatement, query);
query.append(" a.locality = ? ");
preparedStatement.add(locality);
}
// Test with connection number
// addClauseIfRequired(preparedStatement, query);
// query.append(" conn.connectionno = '0603000900' ");
Expand Down Expand Up @@ -436,6 +442,7 @@ public String getConnectionNumberList(String tenantId, String connectionType, Li
addClauseIfRequired(preparedStatement, query);
query.append(" conn.tenantid = ? ");
preparedStatement.add(tenantId);

addClauseIfRequired(preparedStatement, query);
query.append(" conn.connectionno is not null");

Expand Down Expand Up @@ -686,7 +693,7 @@ public String searchLastDemandGenFromDate(String consumerCode, String tenantId,
}

/**
* @param billIds
// * @param billIds
* @param preparedStmtList
* @param builder
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,7 @@ public List<Demand> generateDemand(CalculationReq request, List<Calculation> cal
/**
* Creates or updates Demand
*
* @param requestInfo The RequestInfo of the calculation request
// * @param requestInfo The RequestInfo of the calculation request
* @param calculations The Calculation Objects for which demand has to be
* generated or updated
*/
Expand Down Expand Up @@ -349,7 +349,7 @@ public List<Demand> generateDemandForBillingCycleInBulk(CalculationReq request,
* @param calculations List of Calculation
* @param masterMap Master MDMS Data
* @return Returns list of demands
* @throws IOException
// * @throws IOException
* @throws JsonMappingException
*/
private List<Demand> createDemand(RequestInfo requestInfo, List<Calculation> calculations,
Expand Down Expand Up @@ -1664,7 +1664,7 @@ private void updateRebate(BigDecimal newAmount, DemandDetailAndCollection latest
*
* @param tenantId TenantId for getting master data.
*/
public void generateDemandForTenantId(String tenantId, RequestInfo requestInfo) {
public void generateDemandForTenantId(String tenantId, String locality, RequestInfo requestInfo) {
requestInfo.getUserInfo().setTenantId(tenantId);
Map<String, Object> billingMasterData = calculatorUtils.loadBillingFrequencyMasterData(requestInfo, tenantId);
long taxPeriodFrom = billingMasterData.get("taxPeriodFrom") == null ? 0l
Expand All @@ -1683,7 +1683,12 @@ public void generateDemandForTenantId(String tenantId, RequestInfo requestInfo)
return;
}

generateDemandForULB(billingMasterData, requestInfo, tenantId, taxPeriodFrom, taxPeriodTo);
if(locality!=null && !locality.trim().isEmpty()){
generateDemandForULB(billingMasterData, requestInfo, tenantId, locality, taxPeriodFrom, taxPeriodTo);
}
else{
generateDemandForULB(billingMasterData, requestInfo, tenantId,null, taxPeriodFrom, taxPeriodTo);
}
}

public String generateDemandForConsumerCode(RequestInfo requestInfo, BulkBillCriteria bulkBillCriteria) {
Expand Down Expand Up @@ -1804,7 +1809,7 @@ public String generateDemandForSingle(Map<String, Object> master, SingleDemand s
.filter(p -> taxPeriodFrom.equals(taxPeriods.get(p).getFromDate())).findFirst().getAsInt();
String cone = singleDemand.getConsumercode();
log.info("Billing master data values for non metered connection:: {}", master);
List<WaterDetails> connectionNos = waterCalculatorDao.getConnectionsNoListforsingledemand(tenantId,
List<WaterDetails> connectionNos = waterCalculatorDao.getConnectionsNoListforsingledemand(tenantId,null,
WSCalculationConstant.nonMeterdConnection, taxPeriodFrom, taxPeriodTo, cone);
int bulkSaveDemandCount = configs.getBulkSaveDemandCount() != null ? configs.getBulkSaveDemandCount() : 1;
log.info("Total Connections: {} and batch count: {}", connectionNos.size(), bulkSaveDemandCount);
Expand Down Expand Up @@ -1870,7 +1875,7 @@ public String generateDemandForSingle(Map<String, Object> master, SingleDemand s
"Pushing calculation req to the kafka topic with bulk data of calculationCriteriaList size: {}",
calculationCriteriaList.size());

wsCalculationProducer.push(configs.getCreateDemand(), calculationReq);
wsCalculationProducer.push(configs.getCreateSingleDemand(), calculationReq);
totalRecordsPushedToKafka = totalRecordsPushedToKafka + calculationCriteriaList.size();
calculationCriteriaList.clear();
connectionNosCount = 0;
Expand All @@ -1894,7 +1899,7 @@ public String generateDemandForSingle(Map<String, Object> master, SingleDemand s
"Pushing calculation last req to the kafka topic with bulk data of calculationCriteriaList size: {}",
calculationCriteriaList.size());

wsCalculationProducer.push(configs.getCreateDemand(), calculationReq);
wsCalculationProducer.push(configs.getCreateSingleDemand(), calculationReq);
totalRecordsPushedToKafka = totalRecordsPushedToKafka + calculationCriteriaList.size();
calculationCriteriaList.clear();
connectionNosCount = 0;
Expand Down Expand Up @@ -1982,9 +1987,9 @@ public List<WaterConnection> getConnectionPendingForDemand(RequestInfo requestIn
* @param requestInfo Request Info
* @param tenantId Tenant Id
*/
public void generateDemandForULB(Map<String, Object> master, RequestInfo requestInfo, String tenantId,
public void generateDemandForULB(Map<String, Object> master, RequestInfo requestInfo, String tenantId, String locality,
Long taxPeriodFrom, Long taxPeriodTo) {
log.info("generateDemandForULB:: {} taxPeriodFrom:: {} taxPeriodTo {}", tenantId, taxPeriodFrom, taxPeriodTo);
log.info("generateDemandForULB:: {} Locality:: {} taxPeriodFrom:: {} taxPeriodTo {}", tenantId, locality, taxPeriodFrom, taxPeriodTo);
try {

// ── Step 1: load tax periods from MDMS ───────────────────────────────
Expand Down Expand Up @@ -2018,8 +2023,14 @@ public void generateDemandForULB(Map<String, Object> master, RequestInfo request
// ── Step 3: fetch connections from DB ─────────────────────────────────
List<WaterDetails> connectionNos;
try {
connectionNos = waterCalculatorDao.getConnectionsNoListforsingledemand(tenantId,
WSCalculationConstant.nonMeterdConnection, taxPeriodFrom, taxPeriodTo, cone);
if(locality!=null && !locality.trim().isEmpty()) {
connectionNos = waterCalculatorDao.getConnectionsNoListforsingledemand(tenantId, locality,
WSCalculationConstant.nonMeterdConnection, taxPeriodFrom, taxPeriodTo, cone);
}
else{
connectionNos = waterCalculatorDao.getConnectionsNoListforsingledemand(tenantId, null,
WSCalculationConstant.nonMeterdConnection, taxPeriodFrom, taxPeriodTo, cone);
}
} catch (Exception e) {
log.error("\u274C DB fetch failed for connections (tenant: {}) | {}", tenantId, e.getMessage(), e);
return;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -697,15 +697,18 @@ public void getBillingPeriod(ArrayList<?> mdmsResponse, RequestInfo requestInfo,
/**
* Generate Demand Based on Time (Monthly, Quarterly, Yearly)
*/
public void generateDemandBasedOnTimePeriod(RequestInfo requestInfo, BulkBillCriteria bulkBillCriteria) {
public void generateDemandBasedOnTimePeriod(RequestInfo requestInfo, BulkDemandCriteria bulkDemandCriteria) {
DateTimeFormatter dateTimeFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
LocalDateTime date = LocalDateTime.now();
log.info("Time schedule start for water demand generation on : " + date.format(dateTimeFormatter));
// List<String> tenantIds = wSCalculationDao.getTenantId();
List<String> tenantIds = new ArrayList<>();
String tenat = requestInfo.getMsgId();
List<String> localities = new ArrayList<>();
// String tenat = requestInfo.getMsgId();
String tenant = bulkDemandCriteria.getTenantId();
String locality = bulkDemandCriteria.getLocality();

if (!tenat.contains("pb")) {
if (!tenant.contains("pb")) {
MdmsCriteriaReq mdmsCriteriaReq = calculatorUtil.gettenants(requestInfo);
StringBuilder url = calculatorUtil.getMdmsSearchUrl();
Object res = repository.fetchResult(url, mdmsCriteriaReq);
Expand All @@ -719,24 +722,29 @@ public void generateDemandBasedOnTimePeriod(RequestInfo requestInfo, BulkBillCri
Object mdmsResObj = resMap.get("MdmsRes");
Map<String, Object> mdmsRes = (Map<String, Object>) mdmsResObj;
Object tenantObj = mdmsRes.get("tenant");
Map<String, Object> tenant = (Map<String, Object>) tenantObj;
Object waterSewerageObj = tenant.get("waterSewerage");
Map<String, Object> tenantMap = (Map<String, Object>) tenantObj;
Object waterSewerageObj = tenantMap.get("waterSewerage");
List<Object> waterSewerageList = (List<Object>) waterSewerageObj;
for (Object obj : waterSewerageList) {
if (obj instanceof Map) {
Map<String, Object> waterSewerageMap = (Map<String, Object>) obj;
Object codeObj = waterSewerageMap.get("code");
Object localitiesObj = waterSewerageMap.get("localities");
if (codeObj != null) {
String code = codeObj.toString();
tenantIds.add(code);
}
}
} else if (localitiesObj != null ) {
String local = localitiesObj.toString();
localities.add(local);
}
}
}

}

} else {
tenantIds.add(tenat);
tenantIds.add(tenant);
localities.add(locality);

}
if (tenantIds.isEmpty()) {
Expand Down Expand Up @@ -766,9 +774,15 @@ public void generateDemandBasedOnTimePeriod(RequestInfo requestInfo, BulkBillCri

CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
try {
log.info("\u25B6\uFE0F Demand generation started for tenant: {}", tenantId);
demandService.generateDemandForTenantId(tenantId, tenantRequestInfo);
log.info("\u2705 Demand generation completed for tenant: {}", tenantId);
if(!localities.isEmpty()){
localities.forEach(localty -> {
demandService.generateDemandForTenantId(tenantId, localty, requestInfo);
});
}else {
log.info("\u25B6\uFE0F Demand generation started for tenant: {}", tenantId);
demandService.generateDemandForTenantId(tenantId, null, requestInfo);
log.info("\u2705 Demand generation completed for tenant: {}", tenantId);
}
} catch (Exception e) {
// Catch everything — one tenant failure must NOT block others
log.error("\u274C Demand generation failed for tenant: {} | {}", tenantId, e.getMessage(), e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -88,9 +88,9 @@ public ResponseEntity<DemandResponse> updateDemands(@RequestBody @Valid RequestI
}

@PostMapping("/_jobscheduler")
public void jobscheduler(@Valid @RequestBody BulkBillReq bulkBillReq) {
wSCalculationService.generateDemandBasedOnTimePeriod(bulkBillReq.getRequestInfo(), bulkBillReq.getBulkBillCriteria());
}
public void jobscheduler(@Valid @RequestBody BulkDemandReq bulkDemandReq) {
wSCalculationService.generateDemandBasedOnTimePeriod(bulkDemandReq.getRequestInfo(), bulkDemandReq.getBulkDemandCriteria());
}

@PostMapping("/_singledemand")
public ResponseEntity<Map<String, Object>> singledemandgen(@Valid @RequestBody SingleDemand singledemand) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package org.egov.wscalculation.web.models;

import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.*;
import org.springframework.validation.annotation.Validated;

import java.util.List;

@Validated
@Data
@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
@Builder
@ToString
public class BulkDemandCriteria {

@JsonProperty("tenantIds")
private List<String> tenantIds;

@JsonProperty("offset")
private Long offset;

@JsonProperty("limit")
private Long limit;

@JsonProperty("consumerCode")
private String consumerCode;

@JsonProperty("tenantId")
private String tenantId;

@JsonProperty("locality")
private String locality;

@JsonProperty("localities")
private List<String> localities;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
package org.egov.wscalculation.web.models;

import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.*;
import org.egov.common.contract.request.RequestInfo;
import org.springframework.validation.annotation.Validated;

import javax.validation.constraints.NotNull;

@Validated
@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
@Builder
@ToString
public class BulkDemandReq {

@JsonProperty("RequestInfo")
@NotNull
private RequestInfo requestInfo;

@JsonProperty("BulkDemandCriteria")
private BulkDemandCriteria bulkDemandCriteria;

}
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,9 @@ server.servlet.context-path=/ws-calculator-v2
app.timezone=UTC

spring.datasource.driver-class-name=org.postgresql.Driver
spring.datasource.url=jdbc:postgresql://localhost:5432/pg
spring.datasource.username=postgres
spring.datasource.password=postgres
#Set context root
#spring.datasource.url=jdbc:postgresql://localhost:5432/pg
#spring.datasource.username=postgres
#spring.datasource.password=postgres

#KAFKA topic configuration
egov.meterservice.createmeterconnection=save-ws-meter-v2
Expand All @@ -33,6 +32,7 @@ spring.kafka.producer.value-serializer=org.springframework.kafka.support.seriali
$KAFKA TOPIC DETAILS
#egov.watercalculatorservice.createdemand.topic=ws-generate-demand
egov.watercalculatorservice.createdemand.topic=ws-scheduler-generate-demand-v2
egov.watercalculatorservice.createsingledemand.topic=ws-scheduler-generate-single-demand-v2
egov.watercalculatorservice.billgenerate.topic=ws-bulk-bill-job-scheduler-v2

#Demand Saved And Failed Topic
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,9 @@ public class WSCalculationConfiguration {

@Value("${egov.watercalculatorservice.createdemand.topic}")
private String createDemand;

@Value("${egov.watercalculatorservice.createsingledemand.topic}")
private String createSingleDemand;

@Value("${ws.demand.based.batch.size}")
private Integer batchSize;
Expand Down
Loading