Skip to content

Commit 466138c

Browse files
committed
fix: make it run
1 parent 519d52d commit 466138c

10 files changed

Lines changed: 27 additions & 17 deletions

File tree

database-to-rest/pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
<description>Pipeline from Db to Rest Endpoint</description>
1717
<properties>
1818
<java.version>24</java.version>
19-
<debezium.version>3.0.5.Final</debezium.version>
19+
<debezium.version>3.3.2.Final</debezium.version>
2020
</properties>
2121
<dependencies>
2222
<dependency>

database-to-rest/src/main/java/de/xxx/dbtorest/kafka/KafkaConfig.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,12 +33,14 @@
3333
import org.springframework.kafka.core.ProducerFactory;
3434
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
3535
import org.springframework.kafka.transaction.KafkaTransactionManager;
36+
import org.springframework.transaction.annotation.EnableTransactionManagement;
3637

3738
import java.util.LinkedHashMap;
3839
import java.util.Map;
3940

4041
@Configuration
4142
@EnableKafka
43+
@EnableTransactionManagement
4244
public class KafkaConfig {
4345
private static final Logger log = LoggerFactory.getLogger(KafkaConfig.class);
4446
public static final String ORDERPRODUCT_TOPIC = "orderproduct-topic";

database-to-rest/src/main/java/de/xxx/dbtorest/kafka/KafkaProducer.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,14 +12,14 @@
1212
*/
1313
package de.xxx.dbtorest.kafka;
1414

15-
import com.fasterxml.jackson.databind.ObjectMapper;
1615
import de.xxx.dbtorest.model.DbChangeDto;
1716
import org.apache.kafka.clients.admin.AdminClient;
1817
import org.slf4j.Logger;
1918
import org.slf4j.LoggerFactory;
2019
import org.springframework.kafka.core.KafkaTemplate;
2120
import org.springframework.kafka.support.SendResult;
2221
import org.springframework.stereotype.Component;
22+
import tools.jackson.databind.json.JsonMapper;
2323

2424
import java.util.concurrent.CompletableFuture;
2525
import java.util.concurrent.TimeUnit;
@@ -28,10 +28,10 @@
2828
public class KafkaProducer {
2929
private static final Logger LOGGER = LoggerFactory.getLogger(KafkaProducer.class);
3030
private final KafkaTemplate<String, String> kafkaTemplate;
31-
private final ObjectMapper objectMapper;
31+
private final JsonMapper objectMapper;
3232
private final AdminClient adminClient;
3333

34-
public KafkaProducer(KafkaTemplate<String, String> kafkaTemplate, AdminClient adminClient, ObjectMapper objectMapper) {
34+
public KafkaProducer(KafkaTemplate<String, String> kafkaTemplate, AdminClient adminClient, JsonMapper objectMapper) {
3535
this.adminClient = adminClient;
3636
this.kafkaTemplate = kafkaTemplate;
3737
this.objectMapper = objectMapper;

database-to-rest/src/main/java/de/xxx/dbtorest/sink/DbChangeSinkService.java

Lines changed: 6 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -22,27 +22,25 @@
2222
import org.springframework.http.MediaType;
2323
import org.springframework.stereotype.Service;
2424
import org.springframework.web.client.RestClient;
25+
import tools.jackson.databind.json.JsonMapper;
26+
27+
import java.util.Optional;
2528

2629

2730
@Service
2831
public class DbChangeSinkService {
2932
private static final Logger LOG = LoggerFactory.getLogger(DbChangeSinkService.class);
30-
private final ObjectMapper objectMapper;
33+
private final JsonMapper objectMapper;
3134
private final RestClient restClient = RestClient.create();
3235
@Value("${rest.endpoint.url}")
3336
private String restEndpointUrl;
3437

35-
public DbChangeSinkService(ObjectMapper objectMapper) {
38+
public DbChangeSinkService(JsonMapper objectMapper) {
3639
this.objectMapper = objectMapper;
3740
}
3841

3942
public void handleDbChange(DbChangeDto dbChangeDto) {
40-
Wrapper wrapper;
41-
try {
42-
wrapper = this.objectMapper.readValue(dbChangeDto.value(), Wrapper.class);
43-
} catch (JsonProcessingException e) {
44-
throw new RuntimeException(e);
45-
}
43+
var wrapper = Optional.ofNullable( dbChangeDto.value()).stream().map(value -> this.objectMapper.readValue(value, Wrapper.class)).findFirst().orElse(new Wrapper(null, null));
4644
LOG.info("DbChange received: {}", dbChangeDto.toString());
4745
this.restClient.post().uri(this.restEndpointUrl.trim() + "/rest/orderproduct").contentType(MediaType.APPLICATION_JSON).body(wrapper).retrieve().toBodilessEntity();
4846
}

event-to-file/src/main/java/de/xxx/eventtofile/kafka/KafkaConfig.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,12 +31,14 @@
3131
import org.springframework.kafka.core.ProducerFactory;
3232
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
3333
import org.springframework.kafka.transaction.KafkaTransactionManager;
34+
import org.springframework.transaction.annotation.EnableTransactionManagement;
3435

3536
import java.util.LinkedHashMap;
3637
import java.util.Map;
3738

3839
@Configuration
3940
@EnableKafka
41+
@EnableTransactionManagement
4042
public class KafkaConfig {
4143
private static final Logger log = LoggerFactory.getLogger(KafkaConfig.class);
4244
public static final String FLIGHT_SOURCE_TOPIC = "flight-source-topic";

soap-to-db/pom.xml

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,11 @@
6060
<groupId>org.springframework.boot</groupId>
6161
<artifactId>spring-boot-starter-webservices</artifactId>
6262
</dependency>
63-
63+
<dependency>
64+
<groupId>org.postgresql</groupId>
65+
<artifactId>postgresql</artifactId>
66+
<scope>runtime</scope>
67+
</dependency>
6468
<dependency>
6569
<groupId>org.springframework.boot</groupId>
6670
<artifactId>spring-boot-starter-actuator-test</artifactId>
@@ -86,11 +90,11 @@
8690
<artifactId>spring-boot-starter-webservices-test</artifactId>
8791
<scope>test</scope>
8892
</dependency>
89-
<!--
9093
<dependency>
9194
<groupId>wsdl4j</groupId>
9295
<artifactId>wsdl4j</artifactId>
9396
</dependency>
97+
<!--
9498
<dependency>
9599
<groupId>org.springframework.kafka</groupId>
96100
<artifactId>spring-kafka</artifactId>

soap-to-db/src/main/java/de/xxx/soaptodb/kafka/KafkaConfig.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,12 +34,14 @@
3434
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
3535
import org.springframework.kafka.transaction.KafkaTransactionManager;
3636
import org.springframework.orm.jpa.JpaTransactionManager;
37+
import org.springframework.transaction.annotation.EnableTransactionManagement;
3738

3839
import java.util.LinkedHashMap;
3940
import java.util.Map;
4041

4142
@Configuration
4243
@EnableKafka
44+
@EnableTransactionManagement
4345
public class KafkaConfig {
4446
private static final Logger log = LoggerFactory.getLogger(KafkaConfig.class);
4547
public static final String COUNTRY_TOPIC = "country-topic";

source-sink/pom.xml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -151,6 +151,7 @@
151151
</args>
152152
<compilerPlugins>
153153
<plugin>spring</plugin>
154+
<plugin>jpa</plugin>
154155
</compilerPlugins>
155156
</configuration>
156157
<dependencies>

source-sink/src/main/kotlin/de/xxx/sourcesink/adapter/event/KafkaConfig.kt

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,10 +22,12 @@ import org.springframework.context.annotation.Configuration
2222
import org.springframework.kafka.annotation.EnableKafka
2323
import org.springframework.kafka.core.KafkaTemplate
2424
import org.springframework.kafka.core.ProducerFactory
25+
import org.springframework.transaction.annotation.EnableTransactionManagement
2526

2627

2728
@Configuration
2829
@EnableKafka
30+
@EnableTransactionManagement
2931
class KafkaConfig(val producerFactory: ProducerFactory<String, String>) {
3032
companion object {
3133
@JvmStatic

source-sink/src/main/kotlin/de/xxx/sourcesink/adapter/event/KafkaProducer.kt

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -12,18 +12,17 @@ limitations under the License.
1212
*/
1313
package de.xxx.sourcesink.adapter.event
1414

15-
import com.fasterxml.jackson.databind.ObjectMapper
1615
import de.xxx.sourcesink.domain.model.FlightDto
17-
import org.apache.kafka.clients.admin.AdminClient
1816
import org.slf4j.Logger
1917
import org.slf4j.LoggerFactory
2018
import org.springframework.kafka.core.KafkaTemplate
2119
import org.springframework.stereotype.Component
20+
import tools.jackson.databind.json.JsonMapper
2221
import java.util.concurrent.TimeUnit
2322

2423
@Component
2524
class KafkaProducer(val kafkaTemplate: KafkaTemplate<String, String>,
26-
val objectMapper: ObjectMapper) {
25+
val objectMapper: JsonMapper) {
2726
private val LOGGER: Logger = LoggerFactory.getLogger(KafkaProducer::class.java)
2827

2928
fun sendFlightMsg(flightDto: FlightDto) {

0 commit comments

Comments
 (0)