Commit dface061 authored by Krishnakanth Balla's avatar Krishnakanth Balla

Initial Commit - 15052023

parents
HELP.md
target/
!.mvn/wrapper/maven-wrapper.jar
!**/src/main/**/target/
!**/src/test/**/target/
### STS ###
.apt_generated
.classpath
.factorypath
.project
.settings
.springBeans
.sts4-cache
### IntelliJ IDEA ###
.idea
*.iws
*.iml
*.ipr
### NetBeans ###
/nbproject/private/
/nbbuild/
/dist/
/nbdist/
/.nb-gradle/
build/
!**/src/main/**/build/
!**/src/test/**/build/
### VS Code ###
.vscode/
This diff is collapsed.
@REM ----------------------------------------------------------------------------
@REM Licensed to the Apache Software Foundation (ASF) under one
@REM or more contributor license agreements. See the NOTICE file
@REM distributed with this work for additional information
@REM regarding copyright ownership. The ASF licenses this file
@REM to you under the Apache License, Version 2.0 (the
@REM "License"); you may not use this file except in compliance
@REM with the License. You may obtain a copy of the License at
@REM
@REM https://www.apache.org/licenses/LICENSE-2.0
@REM
@REM Unless required by applicable law or agreed to in writing,
@REM software distributed under the License is distributed on an
@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
@REM KIND, either express or implied. See the License for the
@REM specific language governing permissions and limitations
@REM under the License.
@REM ----------------------------------------------------------------------------
@REM ----------------------------------------------------------------------------
@REM Maven Start Up Batch script
@REM
@REM Required ENV vars:
@REM JAVA_HOME - location of a JDK home dir
@REM
@REM Optional ENV vars
@REM M2_HOME - location of maven2's installed home dir
@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands
@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a keystroke before ending
@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven
@REM e.g. to debug Maven itself, use
@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000
@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files
@REM ----------------------------------------------------------------------------
@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on'
@echo off
@REM set title of command window
title %0
@REM enable echoing by setting MAVEN_BATCH_ECHO to 'on'
@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO%
@REM set %HOME% to equivalent of $HOME
if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%")
@REM Execute a user defined script before this one
if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre
@REM check for pre script, once with legacy .bat ending and once with .cmd ending
if exist "%USERPROFILE%\mavenrc_pre.bat" call "%USERPROFILE%\mavenrc_pre.bat" %*
if exist "%USERPROFILE%\mavenrc_pre.cmd" call "%USERPROFILE%\mavenrc_pre.cmd" %*
:skipRcPre
@setlocal
set ERROR_CODE=0
@REM To isolate internal variables from possible post scripts, we use another setlocal
@setlocal
@REM ==== START VALIDATION ====
if not "%JAVA_HOME%" == "" goto OkJHome
echo.
echo Error: JAVA_HOME not found in your environment. >&2
echo Please set the JAVA_HOME variable in your environment to match the >&2
echo location of your Java installation. >&2
echo.
goto error
:OkJHome
if exist "%JAVA_HOME%\bin\java.exe" goto init
echo.
echo Error: JAVA_HOME is set to an invalid directory. >&2
echo JAVA_HOME = "%JAVA_HOME%" >&2
echo Please set the JAVA_HOME variable in your environment to match the >&2
echo location of your Java installation. >&2
echo.
goto error
@REM ==== END VALIDATION ====
:init
@REM Find the project base dir, i.e. the directory that contains the folder ".mvn".
@REM Fallback to current working directory if not found.
set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR%
IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir
set EXEC_DIR=%CD%
set WDIR=%EXEC_DIR%
:findBaseDir
IF EXIST "%WDIR%"\.mvn goto baseDirFound
cd ..
IF "%WDIR%"=="%CD%" goto baseDirNotFound
set WDIR=%CD%
goto findBaseDir
:baseDirFound
set MAVEN_PROJECTBASEDIR=%WDIR%
cd "%EXEC_DIR%"
goto endDetectBaseDir
:baseDirNotFound
set MAVEN_PROJECTBASEDIR=%EXEC_DIR%
cd "%EXEC_DIR%"
:endDetectBaseDir
IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig
@setlocal EnableExtensions EnableDelayedExpansion
for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a
@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS%
:endReadAdditionalConfig
SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe"
set WRAPPER_JAR="%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.jar"
set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain
set DOWNLOAD_URL="https://repo.maven.apache.org/maven2/org/apache/maven/wrapper/maven-wrapper/3.1.0/maven-wrapper-3.1.0.jar"
FOR /F "usebackq tokens=1,2 delims==" %%A IN ("%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.properties") DO (
IF "%%A"=="wrapperUrl" SET DOWNLOAD_URL=%%B
)
@REM Extension to allow automatically downloading the maven-wrapper.jar from Maven-central
@REM This allows using the maven wrapper in projects that prohibit checking in binary data.
if exist %WRAPPER_JAR% (
if "%MVNW_VERBOSE%" == "true" (
echo Found %WRAPPER_JAR%
)
) else (
if not "%MVNW_REPOURL%" == "" (
SET DOWNLOAD_URL="%MVNW_REPOURL%/org/apache/maven/wrapper/maven-wrapper/3.1.0/maven-wrapper-3.1.0.jar"
)
if "%MVNW_VERBOSE%" == "true" (
echo Couldn't find %WRAPPER_JAR%, downloading it ...
echo Downloading from: %DOWNLOAD_URL%
)
powershell -Command "&{"^
"$webclient = new-object System.Net.WebClient;"^
"if (-not ([string]::IsNullOrEmpty('%MVNW_USERNAME%') -and [string]::IsNullOrEmpty('%MVNW_PASSWORD%'))) {"^
"$webclient.Credentials = new-object System.Net.NetworkCredential('%MVNW_USERNAME%', '%MVNW_PASSWORD%');"^
"}"^
"[Net.ServicePointManager]::SecurityProtocol = [Net.SecurityProtocolType]::Tls12; $webclient.DownloadFile('%DOWNLOAD_URL%', '%WRAPPER_JAR%')"^
"}"
if "%MVNW_VERBOSE%" == "true" (
echo Finished downloading %WRAPPER_JAR%
)
)
@REM End of extension
@REM Provide a "standardized" way to retrieve the CLI args that will
@REM work with both Windows and non-Windows executions.
set MAVEN_CMD_LINE_ARGS=%*
%MAVEN_JAVA_EXE% ^
%JVM_CONFIG_MAVEN_PROPS% ^
%MAVEN_OPTS% ^
%MAVEN_DEBUG_OPTS% ^
-classpath %WRAPPER_JAR% ^
"-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" ^
%WRAPPER_LAUNCHER% %MAVEN_CONFIG% %*
if ERRORLEVEL 1 goto error
goto end
:error
set ERROR_CODE=1
:end
@endlocal & set ERROR_CODE=%ERROR_CODE%
if not "%MAVEN_SKIP_RC%"=="" goto skipRcPost
@REM check for post script, once with legacy .bat ending and once with .cmd ending
if exist "%USERPROFILE%\mavenrc_post.bat" call "%USERPROFILE%\mavenrc_post.bat"
if exist "%USERPROFILE%\mavenrc_post.cmd" call "%USERPROFILE%\mavenrc_post.cmd"
:skipRcPost
@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on'
if "%MAVEN_BATCH_PAUSE%"=="on" pause
if "%MAVEN_TERMINATE_CMD%"=="on" exit %ERROR_CODE%
cmd /C exit /B %ERROR_CODE%
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.5.7</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<groupId>com.mongo.reactivie</groupId>
<artifactId>mongoreactive</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>mongoreactive</name>
<description>Demo project for Spring Boot</description>
<properties>
<java.version>11</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-mongodb-reactive</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-mongodb</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor.kafka</groupId>
<artifactId>reactor-kafka</artifactId>
<version>1.3.5</version>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-collections4</artifactId>
<version>4.4</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- https://mvnrepository.com/artifact/io.springfox/springfox-swagger-ui -->
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger-ui</artifactId>
<version>2.9.2</version>
</dependency>
<!-- https://mvnrepository.com/artifact/io.springfox/springfox-swagger2 -->
<dependency>
<groupId>io.springfox</groupId>
<artifactId>springfox-swagger2</artifactId>
<version>2.9.2</version>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>
</project>
package com.mongo.reactivie;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.reactive.config.EnableWebFlux;
import springfox.documentation.builders.ApiInfoBuilder;
import springfox.documentation.builders.PathSelectors;
import springfox.documentation.builders.RequestHandlerSelectors;
import springfox.documentation.service.ApiInfo;
import springfox.documentation.spi.DocumentationType;
import springfox.documentation.spring.web.plugins.Docket;
@SpringBootApplication
@EnableWebFlux
public class MongoreactiveApplication {
public static void main(String[] args) {
SpringApplication.run(MongoreactiveApplication.class, args);
}
}
package com.mongo.reactivie.config;
import com.mongo.reactivie.model.JobPost;
import java.util.Collections;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.reactive.ReactiveKafkaConsumerTemplate;
import reactor.kafka.receiver.ReceiverOptions;
@Configuration
public class KafkaConsumerConfig {
@Bean
public ReceiverOptions<String, JobPost> kafkaReceiverOptions(@Value(value = "${FAKE_CONSUMER_DTO_TOPIC}") String topic, KafkaProperties kafkaProperties) {
ReceiverOptions<String, JobPost> basicReceiverOptions = ReceiverOptions.create(kafkaProperties.buildConsumerProperties());
return basicReceiverOptions.subscription(Collections.singletonList(topic));
}
@Bean
public ReactiveKafkaConsumerTemplate<String, JobPost> reactiveKafkaConsumerTemplate(ReceiverOptions<String, JobPost> kafkaReceiverOptions) {
return new ReactiveKafkaConsumerTemplate<String, JobPost>(kafkaReceiverOptions);
}
}
package com.mongo.reactivie.config;
import java.util.Map;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.reactive.ReactiveKafkaProducerTemplate;
import reactor.kafka.sender.SenderOptions;
@Configuration
public class KafkaProducerConfig {
@Bean
public ReactiveKafkaProducerTemplate<String, String> reactiveKafkaProducerTemplate(
KafkaProperties kafkaProperties) {
Map<String, Object> props = kafkaProperties.buildProducerProperties();
return new ReactiveKafkaProducerTemplate<String, String>(SenderOptions.create(props));
}
}
package com.mongo.reactivie.consumer;
import com.mongo.reactivie.model.JobPost;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.CommandLineRunner;
import org.springframework.kafka.core.reactive.ReactiveKafkaConsumerTemplate;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
@Service
public class ReactiveConsumerService implements CommandLineRunner {
Logger log = LoggerFactory.getLogger(ReactiveConsumerService.class);
private final ReactiveKafkaConsumerTemplate<String, JobPost> reactiveKafkaConsumerTemplate;
public ReactiveConsumerService(ReactiveKafkaConsumerTemplate<String, JobPost> reactiveKafkaConsumerTemplate) {
this.reactiveKafkaConsumerTemplate = reactiveKafkaConsumerTemplate;
}
private Flux<JobPost> consumeJobPost() {
return reactiveKafkaConsumerTemplate
.receiveAutoAck()
.doOnNext(consumerRecord -> log.info("received key={}, value={} from topic={}, offset={}",
consumerRecord.key(),
consumerRecord.value(),
consumerRecord.topic(),
consumerRecord.offset())
)
.map(ConsumerRecord::value)
.doOnNext(jobPost -> log.info("successfully consumed {}={}", JobPost.class.getSimpleName(), jobPost))
.doOnError(throwable -> log.error("something bad happened while consuming : {}", throwable.getMessage()));
}
@Override
public void run(String... args) {
// we have to trigger consumption
consumeJobPost().subscribe();
}
}
package com.mongo.reactivie.controller;
import com.mongo.reactivie.model.JobPost;
import com.mongo.reactivie.model.JobPostCompanyDetails;
import com.mongo.reactivie.service.JobService;
import java.util.List;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@RestController
public class JobPostController {
@Autowired
private JobService jobService;
@GetMapping("/posts")
@ResponseStatus(HttpStatus.OK)
public Flux<JobPost> getAllPosts() {
return jobService.findAll();
}
@GetMapping("/posts/{id}")
@ResponseStatus(HttpStatus.OK)
public Mono<JobPost> getPostById(@PathVariable("id") String id) {
return jobService.findById(id);
}
@PostMapping("/posts")
@ResponseStatus(HttpStatus.CREATED)
public Mono<JobPost> createPost(@RequestBody JobPost jobPost) {
return jobService.save(jobPost);
}
@PostMapping("/posts/saveall")
@ResponseStatus(HttpStatus.CREATED)
public Flux<JobPost> saveAll(@RequestBody List<JobPost> jobPosts) {
return jobService.saveAll(jobPosts);
}
@PutMapping("/posts/{id}")
@ResponseStatus(HttpStatus.OK)
public Mono<JobPost> updatePost(@PathVariable("id") String id, @RequestBody JobPost jobPost) {
return jobService.update(id, jobPost);
}
@DeleteMapping("/posts/{id}")
@ResponseStatus(HttpStatus.NO_CONTENT)
public Mono<Void> deletePost(@PathVariable("id") String id) {
return jobService.deleteById(id);
}
@GetMapping("/publishMessage/{id}")
@ResponseStatus(HttpStatus.OK)
public Mono<String> publishMessage(@PathVariable("id") String id) {
return jobService.publishMessage(id);
}
@GetMapping("/posts/getJobPostCompanyDetails")
@ResponseStatus(HttpStatus.OK)
public Mono<List<JobPostCompanyDetails>> getJobPostCompanyDetails() {
return Mono.just(jobService.getJobPostCompanyDetails());
}
}
package com.mongo.reactivie.model;
import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;
@Document(collection = "Company")
public class Company {
@Id
private String id;
private String companyName;
private String location;
public String getId() {
return id;
}
public void setId(String id) {
this.id = id;
}
public String getCompanyName() {
return companyName;
}
public void setCompanyName(String companyName) {
this.companyName = companyName;
}
public String getLocation() {
return location;
}
public void setLocation(String location) {
this.location = location;
}
}
package com.mongo.reactivie.model;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonRootName;
import java.util.List;
import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;
@Document(collection = "JobPost")
public class JobPost {
@Id
private String id;
private String profile;
private String desc;
private Integer exp;
private List<String> techs;
private String companyId;
public String getId() {
return id;
}
public void setId(String id) {
this.id = id;
}
public String getProfile() {
return profile;
}
public void setProfile(String profile) {
this.profile = profile;
}
public String getDesc() {
return desc;
}
public void setDesc(String desc) {
this.desc = desc;
}
public Integer getExp() {
return exp;
}
public void setExp(Integer exp) {
this.exp = exp;
}
public List<String> getTechs() {
return techs;
}
public void setTechs(List<String> techs) {
this.techs = techs;
}
public String getCompanyId() {
return companyId;
}
public void setCompanyId(String companyId) {
this.companyId = companyId;
}
@Override
public String toString() {
return "JobPost{" +
"id='" + id + '\'' +
", profile='" + profile + '\'' +
", desc='" + desc + '\'' +
", exp=" + exp +
", techs=" + techs +
'}';
}
}
package com.mongo.reactivie.model;
import java.util.List;
public class JobPostCompanyDetails {
private String profile;
private String desc;
private Integer exp;
private List<String> techs;
private List<Company> companyDetails;
public String getProfile() {
return profile;
}
public void setProfile(String profile) {
this.profile = profile;
}
public String getDesc() {
return desc;
}
public void setDesc(String desc) {
this.desc = desc;
}
public Integer getExp() {
return exp;
}
public void setExp(Integer exp) {
this.exp = exp;
}
public List<String> getTechs() {
return techs;
}
public void setTechs(List<String> techs) {
this.techs = techs;
}
public List<Company> getCompanyDetails() {
return companyDetails;
}
public void setCompanyDetails(List<Company> companyDetails) {
this.companyDetails = companyDetails;
}
}
package com.mongo.reactivie.model.deserializer;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonMappingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.mongo.reactivie.model.JobPost;
import java.util.Optional;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.serialization.Deserializer;
import org.springframework.stereotype.Component;
@Component
public class JobPostDeserializer implements Deserializer<JobPost> {
@Override
public JobPost deserialize(String arg0, byte[] arg1) {
ObjectMapper mapper = new ObjectMapper();
JobPost jobPost = new JobPost();
try {
JsonNode orderNode = mapper.readTree(arg1);
if (Optional.ofNullable(orderNode).isPresent()) {
constructJobPost(jobPost, orderNode);
}
} catch (Exception e) {
}
return jobPost;
}
private void constructJobPost(JobPost jobPost, JsonNode jsonNode)
throws JsonMappingException, JsonProcessingException {
jobPost.setId(getTextValue(jsonNode, "id"));
jobPost.setDesc(getTextValue(jsonNode, "desc"));
jobPost.setExp((getIntValue(jsonNode, "exp")));
jobPost.setProfile(getTextValue(jsonNode, "profile"));
}
private String getTextValue(JsonNode jsonNode, String key) {
return Optional.ofNullable(jsonNode.get(key)).map(JsonNode::textValue).orElse(null);
}
private Integer getIntValue(JsonNode jsonNode, String key) {
return Optional.ofNullable(jsonNode.get(key)).map(JsonNode::intValue).orElse(null);
}
@Override
public JobPost deserialize(String arg0, Headers arg1, byte[] arg2) {
return this.deserialize(arg0, arg2);
}
}
package com.mongo.reactivie.repository;
import com.mongo.reactivie.model.JobPostCompanyDetails;
import com.mongodb.client.AggregateIterable;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoCollection;
import com.mongodb.client.MongoDatabase;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import org.bson.Document;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.mongodb.core.convert.MongoConverter;
import org.springframework.stereotype.Component;
@Component
public class JobPostDAOImpl {
@Autowired
private MongoClient mongoClient;
@Autowired
MongoConverter converter;
public List<JobPostCompanyDetails> getJobPostCompanyDetails() {
final List<JobPostCompanyDetails> posts = new ArrayList<>();
MongoDatabase database = mongoClient.getDatabase("Job");
MongoCollection<Document> collection = database.getCollection("JobPost");
AggregateIterable<Document> result = collection.aggregate(Arrays.asList(new Document("$lookup",
new Document("from", "Company")
.append("localField", "companyId")
.append("foreignField", "_id")
.append("as", "companyDetails"))));
result.forEach(doc -> posts.add(converter.read(JobPostCompanyDetails.class,doc)));
return posts;
}
}
\ No newline at end of file
package com.mongo.reactivie.repository;
import com.mongo.reactivie.model.JobPost;
import org.springframework.data.mongodb.repository.ReactiveMongoRepository;
import org.springframework.stereotype.Repository;
@Repository
public interface JobPostRepository extends ReactiveMongoRepository<JobPost, String> {
}
package com.mongo.reactivie.service;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.mongo.reactivie.model.JobPost;
import com.mongo.reactivie.model.JobPostCompanyDetails;
import com.mongo.reactivie.repository.JobPostDAOImpl;
import com.mongo.reactivie.repository.JobPostRepository;
import java.util.List;
import java.util.Optional;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@Service
public class JobService {
@Autowired
private JobPostRepository jobPostRepository;
@Autowired
private MessagePublishService messagePublishService;
@Autowired
private JobPostDAOImpl jobPostDAOImpl;
public Flux<JobPost> findAll() {
return jobPostRepository.findAll();
}
public Mono<JobPost> findById(String id) {
return jobPostRepository.findById(id);
}
public Mono<JobPost> save(JobPost jobPost) {
return jobPostRepository.save(jobPost);
}
public Flux<JobPost> saveAll(List<JobPost> jobPosts) {
return jobPostRepository.saveAll(jobPosts);
}
public Mono<JobPost> update(String id, JobPost jobPost) {
return jobPostRepository.findById(id).map(Optional::of).defaultIfEmpty(Optional.empty())
.flatMap(optionalJobPost -> {
if (optionalJobPost.isPresent()) {
jobPost.setId(id);
return jobPostRepository.save(jobPost);
}
return Mono.empty();
});
}
public Mono<Void> deleteById(String id) {
return jobPostRepository.deleteById(id);
}
public Mono<String> publishMessage(String id) {
return jobPostRepository.findById(id)
.flatMap(post -> messagePublishService.publishMessage(convertToJson(post)));
}
public List<JobPostCompanyDetails> getJobPostCompanyDetails(){
return jobPostDAOImpl.getJobPostCompanyDetails();
}
private String convertToJson(JobPost posts) {
ObjectMapper mapper = new ObjectMapper();
String json = " ";
try {
json = mapper.writeValueAsString(posts);
} catch (JsonProcessingException e) {
e.printStackTrace();
}
return json;
}
}
package com.mongo.reactivie.service;
import java.util.Map;
import org.apache.commons.collections4.map.HashedMap;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.core.reactive.ReactiveKafkaProducerTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Service;
import reactor.core.publisher.Mono;
@Service
public class MessagePublishService {
private final ReactiveKafkaProducerTemplate<String, String> reactiveKafkaProducerTemplate;
public MessagePublishService(
ReactiveKafkaProducerTemplate<String, String> reactiveKafkaProducerTemplate) {
this.reactiveKafkaProducerTemplate = reactiveKafkaProducerTemplate;
}
public Mono<String> publishMessage(String jsonString) {
ProducerRecord<String, String> record = new ProducerRecord<>(
"TOPIC", jsonString);
Map<String, String> headers = new HashedMap<>();
headers.forEach((key, value) -> record.headers().add(key, value.getBytes()));
reactiveKafkaProducerTemplate.send(record).subscribe();
return Mono.just("Message published successfully");
}
}
spring.data.mongodb.uri=mongodb+srv://root:root@cluster0.sr2lgkq.mongodb.net/?retryWrites=true&w=majority
spring.data.mongodb.database = Job
spring.kafka.bootstrap-servers=localhost:9092
# producer
spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer
spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer
# consumer
spring.kafka.consumer.group-id=reactivekafkaconsumerandproducer
spring.kafka.consumer.auto-offset-reset=latest
spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer
spring.kafka.consumer.value-deserializer=com.mongo.reactivie.model.deserializer.JobPostDeserializer
# json deserializer config
spring.kafka.properties.spring.json.trusted.packages=*
spring.kafka.consumer.properties.spring.json.use.type.headers=false
FAKE_CONSUMER_DTO_TOPIC=TOPIC
package com.mongo.reactivie;
import org.junit.jupiter.api.Test;
import org.springframework.boot.test.context.SpringBootTest;
@SpringBootTest
class MongoreactiveApplicationTests {
@Test
void contextLoads() {
}
}
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment