Commit cd6e981e authored by Potharaju Peddi's avatar Potharaju Peddi

for webflux reactive kafka with mongoDB

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/
# Licensed to the Apache Software Foundation (ASF) under one
# or more contributor license agreements. See the NOTICE file
# distributed with this work for additional information
# regarding copyright ownership. The ASF licenses this file
# to you under the Apache License, Version 2.0 (the
# "License"); you may not use this file except in compliance
# with the License. You may obtain a copy of the License at
#
# https://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing,
# software distributed under the License is distributed on an
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.8.7/apache-maven-3.8.7-bin.zip
wrapperUrl=https://repo.maven.apache.org/maven2/org/apache/maven/wrapper/maven-wrapper/3.1.1/maven-wrapper-3.1.1.jar
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.7.11-SNAPSHOT</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<groupId>com.WebFluxReactiveMongo</groupId>
<artifactId>spring-reactive-mango</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>spring-reactive-mango</name>
<description>Demo project for Spring Boot</description>
<properties>
<java.version>1.8</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-webflux</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor.kafka</groupId>
<artifactId>reactor-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
<configuration>
<excludes>
<exclude>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</exclude>
</excludes>
</configuration>
</plugin>
</plugins>
</build>
<repositories>
<repository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</repository>
<repository>
<id>spring-snapshots</id>
<name>Spring Snapshots</name>
<url>https://repo.spring.io/snapshot</url>
<releases>
<enabled>false</enabled>
</releases>
</repository>
</repositories>
<pluginRepositories>
<pluginRepository>
<id>spring-milestones</id>
<name>Spring Milestones</name>
<url>https://repo.spring.io/milestone</url>
<snapshots>
<enabled>false</enabled>
</snapshots>
</pluginRepository>
<pluginRepository>
<id>spring-snapshots</id>
<name>Spring Snapshots</name>
<url>https://repo.spring.io/snapshot</url>
<releases>
<enabled>false</enabled>
</releases>
</pluginRepository>
</pluginRepositories>
</project>
package com.WebFluxReactiveMongo.Reactor.Controller;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import com.WebFluxReactiveMongo.Reactor.Service.StudentService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@RestController
@RequestMapping("/students")
@Slf4j
public class StudentController {
@Autowired
private StudentService studentService;
@GetMapping
public Flux<StudentDto> getStudents(){
return studentService.getStudents();
}
/*@GetMapping
public Flux<Student> getStudentsTemp(){
log.info("inside controller=========");
return studentService.getStudentsWithTemp();
}*/
@GetMapping("/{id}")
public Mono<StudentDto> studentById(@PathVariable String id){
return studentService.getStudentById(id);
}
@GetMapping("/student-fee-range")
public Flux<StudentDto> getStudentFeeRange(@RequestParam("min") double min, @RequestParam("max") double max){
return studentService.getStudentInFee(min,max);
}
@PostMapping
public Mono<StudentDto> saveStudent(@RequestBody Mono<StudentDto> studentDtoMono){
return studentService.saveStudent(studentDtoMono);
}
@PutMapping("/update/{id}")
public Mono<StudentDto> updateStudent(@RequestBody Mono<StudentDto> studentDtoMono, @PathVariable String id){
return studentService.updateStudent(studentDtoMono,id);
}
@DeleteMapping("/{id}")
public Mono<Void> deleteStudent(@PathVariable String id){
return studentService.deleteStudent(id);
}
}
package com.WebFluxReactiveMongo.Reactor.Dto;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor
@AllArgsConstructor
public class StudentDto {
private String id;
private String name;
private String college;
private double fee;
}
package com.WebFluxReactiveMongo.Reactor.Entity;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;
@Data
@NoArgsConstructor
@AllArgsConstructor
@Document(collection = "student_tbl")
public class Student {
@Id
private String id;
private String name;
private String college;
private double fee;
}
package com.WebFluxReactiveMongo.Reactor.MongoDB;
import com.mongodb.DBObject;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.mongodb.core.ReactiveMongoTemplate;
public class MongoExamples {
}
package com.WebFluxReactiveMongo.Reactor.Repository;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import com.WebFluxReactiveMongo.Reactor.Entity.Student;
import org.springframework.data.domain.Range;
import org.springframework.data.mongodb.repository.ReactiveMongoRepository;
import org.springframework.stereotype.Repository;
import reactor.core.publisher.Flux;
@Repository
public interface StudentRepository extends ReactiveMongoRepository<Student,String> {
Flux<StudentDto> findByFeeBetween(Range<Double> feeRange);
Flux<StudentDto> findByIdBetween(Range<Integer> idRange);
}
package com.WebFluxReactiveMongo.Reactor.Service;
import com.WebFluxReactiveMongo.Reactor.Entity.Student;
import com.WebFluxReactiveMongo.Reactor.Util.ApplicaitonUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.Range;
import org.springframework.data.domain.Sort;
import org.springframework.data.mongodb.core.ReactiveMongoTemplate;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.stereotype.Service;
import com.WebFluxReactiveMongo.Reactor.Repository.StudentRepository;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@Service
@Slf4j
public class StudentService {
@Autowired
private StudentRepository studentRepository;
@Autowired
private ReactiveMongoTemplate reactiveMongoTemplate;
public StudentService(ReactiveMongoTemplate reactiveMongoTemplate){
this.reactiveMongoTemplate = reactiveMongoTemplate;
}
public Flux<Student> getStudentsWithTemp(){
Query query = new Query()
.addCriteria(Criteria.where("name").is("raju"))
.with(Sort.by(Sort.Order.desc("fee")))
.limit(3);
log.info("query======="+query);
return reactiveMongoTemplate.find(query, Student.class);
}
public Flux<StudentDto> getStudents(){
return studentRepository.findAll().map(ApplicaitonUtil::entityToDto);
}
public Mono<StudentDto> getStudentById(String id){
return studentRepository.findById(id).map(ApplicaitonUtil::entityToDto);
}
public Flux<StudentDto> getStudentInFee(double min, double max){
return studentRepository.findByFeeBetween(Range.closed(min,max));
}
public Flux<StudentDto> getStudentInId(int min,int max){
return studentRepository.findByIdBetween(Range.closed(min,max));
}
public Mono<StudentDto> saveStudent(Mono<StudentDto> studentDtoMono){
return studentDtoMono.map(ApplicaitonUtil::dtoToEntity).
flatMap(studentRepository::insert).
map(ApplicaitonUtil::entityToDto);
}
public Mono<StudentDto> updateStudent(Mono<StudentDto> studentDtoMono, String id){
return studentRepository.findById(id).
flatMap(p->studentDtoMono.map(ApplicaitonUtil::dtoToEntity)
.doOnNext(e->e.setId(id))) .flatMap(studentRepository::save)
.map(ApplicaitonUtil::entityToDto);
}
public Mono<Void> deleteStudent(String id){
return studentRepository.deleteById(id);
}
}
package com.WebFluxReactiveMongo.Reactor;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class SpringReactiveMangoApplication {
public static void main(String[] args) {
SpringApplication.run(SpringReactiveMangoApplication.class, args);
}
}
package com.WebFluxReactiveMongo.Reactor.Util;
import com.WebFluxReactiveMongo.Reactor.Entity.Student;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import org.springframework.beans.BeanUtils;
public class ApplicaitonUtil {
public static StudentDto entityToDto(Student student){
StudentDto studentDto = new StudentDto();
BeanUtils.copyProperties(student, studentDto);
return studentDto;
}
public static Student dtoToEntity(StudentDto studentDto){
Student student = new Student();
BeanUtils.copyProperties(studentDto, student);
return student;
}
}
package com.WebFluxReactiveMongo.Reactor.kafka;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import com.WebFluxReactiveMongo.Reactor.Service.StudentService;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Autowired;
import reactor.core.Disposable;
import reactor.core.publisher.Flux;
import reactor.kafka.receiver.ReceiverOptions;
import reactor.kafka.receiver.ReceiverRecord;
import reactor.kafka.receiver.KafkaReceiver;
@Slf4j
public class SampleConsumer {
@Autowired
private StudentService studentService;
private static final String BOOTSTRAP_SERVERS = "localhost:9092";
private static final String TOPIC = "TestTopic";
private final ReceiverOptions<String, JsonNode> receiverOptions;
public SampleConsumer(String bootstrapServers) {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ConsumerConfig.CLIENT_ID_CONFIG, "sample-consumer");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "sample-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
receiverOptions = ReceiverOptions.create(props);
}
public Disposable consumeMessages(String topic, CountDownLatch latch) {
ReceiverOptions<String, JsonNode> options = receiverOptions.subscription(Collections.singleton(topic))
.addAssignListener(partitions -> log.debug("onPartitionsAssigned {}", partitions))
.addRevokeListener(partitions -> log.debug("onPartitionsRevoked {}", partitions));
Flux<ReceiverRecord<String, JsonNode>> kafkaFlux = KafkaReceiver.create(options).receive();
return kafkaFlux
.subscribe(r -> {
log.info("Received message {} ", r);
r.receiverOffset().acknowledge();
latch.countDown();
});
}
public static void main(String[] args) throws Exception {
int count = 1;
CountDownLatch latch = new CountDownLatch(count);
SampleConsumer consumer = new SampleConsumer(BOOTSTRAP_SERVERS);
Disposable disposable = consumer.consumeMessages(TOPIC, latch);
latch.await(10, TimeUnit.SECONDS);
disposable.dispose();
}
}
package com.WebFluxReactiveMongo.Reactor.kafka;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import com.WebFluxReactiveMongo.Reactor.Repository.StudentRepository;
import com.WebFluxReactiveMongo.Reactor.Service.StudentService;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.support.serializer.JsonSerializer;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.kafka.sender.KafkaSender;
import reactor.kafka.sender.SenderOptions;
import reactor.kafka.sender.SenderRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.common.serialization.IntegerSerializer;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
@Slf4j
public class SampleProducer {
private static final String BOOTSTRAP_SERVERS = "localhost:9092";
private static final String TOPIC = "TestTopic1";
@Autowired
private StudentService studentService;
@Autowired
private StudentRepository studentRepository;
private final KafkaSender<String, JsonNode> sender;
public SampleProducer(String bootstrapServers) {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.CLIENT_ID_CONFIG, "sample-producer");
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
SenderOptions<String, JsonNode> senderOptions = SenderOptions.create(props);
sender = KafkaSender.create(senderOptions);
}
public void sendMessages(String topic, int count, CountDownLatch latch) throws InterruptedException {
ObjectMapper objectMapper = new ObjectMapper();
//Flux<StudentDto> studentDto = Flux.just(new StudentDto("106","AADHI","DIYA",50000),
//new StudentDto("107","VASU","KIYA",45000));
sender.<Integer>send(Flux.range(1, 1)
//.map(i -> SenderRecord.create(new ProducerRecord<>(topic, i, "Message_" + i), i)))
.map(i -> {
JsonNode jsonValue = null;
try {
Mono<StudentDto> studentDtoMono=Mono.just(new StudentDto("5855","Kafka Producer message","Nisum School",60000));
String value = objectMapper.writeValueAsString(studentDtoMono);
jsonValue = objectMapper.readTree(value);
} catch (JsonProcessingException e) {
e.printStackTrace();
}
log.info("Message: {}, sent successfully", jsonValue);
return SenderRecord.create(new ProducerRecord<>(topic, "Key_" + UUID.randomUUID(), jsonValue), i);
}))
.doOnError(e -> log.error("Send failed", e))
.doOnNext(r -> log.debug("Message #{} result: {}", r.correlationMetadata(), r.recordMetadata()))
.subscribe();
latch.countDown();
}
public void close() {
sender.close();
}
public static void main(String[] args) throws Exception {
int count = 20;
CountDownLatch latch = new CountDownLatch(count);
SampleProducer producer = new SampleProducer(BOOTSTRAP_SERVERS);
producer.sendMessages(TOPIC, count, latch);
latch.await(10, TimeUnit.SECONDS);
producer.close();
}
}
spring:
data:
mongodb:
database: StudentsDb
host: localhost
port: 27017
\ No newline at end of file
package com.WebFluxReactiveMongo.Reactor;
import com.WebFluxReactiveMongo.Reactor.Controller.StudentController;
import com.WebFluxReactiveMongo.Reactor.Dto.StudentDto;
import com.WebFluxReactiveMongo.Reactor.Service.StudentService;
import org.junit.Test;
//import org.junit.jupiter.api.Test;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.autoconfigure.web.reactive.WebFluxTest;
import org.springframework.boot.test.mock.mockito.MockBean;
import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.test.web.reactive.server.WebTestClient;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import static org.mockito.BDDMockito.given;
import static org.mockito.Mockito.when;
import static org.mockito.ArgumentMatchers.any;
@RunWith(SpringRunner.class)
@WebFluxTest(StudentController.class)
public class SpringReactiveMangoApplicationTests {
@Autowired
private WebTestClient webTestClient;
@MockBean
private StudentService studentService;
@Test
public void addStudentTest(){
Mono<StudentDto> studentDtoMono = Mono.just(new StudentDto("101","RAJU","CJJC",7000));
when(studentService.saveStudent(studentDtoMono)).thenReturn(studentDtoMono);
webTestClient.post().uri("/students").body(Mono.just(studentDtoMono), StudentDto.class)
.exchange().expectStatus().isOk(); //200
}
@Test
public void getStudentsTest(){
Flux<StudentDto> studentDtoFlux = Flux.just(new StudentDto("102","RAVI","GDC",1500),
new StudentDto("103","RAMU","KRDC",5500));
when(studentService.getStudents()).thenReturn(studentDtoFlux);
Flux<StudentDto> responseBody = webTestClient.get().uri("/students")
.exchange().expectStatus().isOk()
.returnResult(StudentDto.class)
.getResponseBody();
StepVerifier.create(responseBody)
.expectSubscription()
.expectNext(new StudentDto("102","RAVI","GDC",1500))
.expectNext(new StudentDto("103","RAMU","KRDC",5500))
.verifyComplete();
}
@Test
public void getStudentTest(){
Mono<StudentDto> studentDtoMono = Mono.just(new StudentDto("102","RAVI","GDC",1500));
when(studentService.getStudentById(any())).thenReturn(studentDtoMono);
Flux<StudentDto> responseBody = webTestClient.get().uri("/students/102")
.exchange().expectStatus().isOk()
.returnResult(StudentDto.class)
.getResponseBody();
StepVerifier.create(responseBody)
.expectSubscription()
.expectNextMatches(p->p.getName().equals("RAVI"))
.verifyComplete();
}
@Test
public void updateStudentTest(){
Mono<StudentDto> studentDtoMono = Mono.just(new StudentDto("102","RAVI","GDC",1500));
when(studentService.updateStudent(studentDtoMono,"102")).thenReturn(studentDtoMono);
webTestClient.put().uri("/students/update/102")
.exchange().expectStatus().isOk();
}
@Test
public void deleteStudentTest(){
given(studentService.deleteStudent("102")).willReturn(Mono.empty());
webTestClient.delete().uri("/students/102")
.exchange().expectStatus().isOk();
}
}
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