- Agile 2
- Alta disponibilidad 1
- Alternativas cloud 1
- Aop 1
- Arquitectura 3
- Arquitectura distribuida 2
- Automatizacion 3
- Azure devops 1
- Base de datos 1
- Buenas practicas 19
- Cloud 1
- Colas 7
- Competing consumers 1
- Convenciones 11
- Copilot 1
- Diseno 6
- Docker 2
- Docker compose 1
- Documentacion 1
- Eda 11
- Equipos 1
- Escalabilidad 1
- Flujo de negocio 1
- Flujo de trabajo 3
- Flyway 1
- Git 4
- Gradle 3
- Herramientas digitales 1
- Ia 1
- Iam 1
- Infraestructura 2
- Java 14
- Jerarquia tecnica 1
- Jpa 1
- Jsonb 1
- Kafka 7
- Kubernetes 1
- Liderazgo en software 1
- Lineamientos 1
- Log 1
- Logging 3
- Microservicios 3
- Mongodb 1
- Monitoreo 1
- Nosql 3
- Observabilidad 4
- Open source 1
- Plugins 3
- Postgresql 1
- Privacidad 1
- Programacion funcional 1
- Programacion reactiva 4
- Rabbitmq 6
- Rotacion de talento 1
- Saga 2
- Scrum 2
- Security 1
- Seguridad 1
- Self hosting 1
- Sistemas legados 1
- Spring boot 3
- Spring mvc 2
- Sql 3
- Streams 1
- Threadlocal 1
- Trazabilidad 2
- Versionado 2
- Web 1
- Webflux 2
- Websockets 1
- Zero trust 1
Webflux
2 artículos
Respuestas API Consistentes: Un Wrapper Transversal con Spring WebFlux y `WebFilter`
- Mauricio ECR
- Snippets
- 30 Aug, 2025
En entornos modernos de microservicios, la consistencia en las respuestas de una API es más que una cuestión de estética: es un factor crítico para la mantenibilidad, observabilidad y experiencia del
Respuestas API Consistentes: Un Wrapper Transversal con Spring WebFlux y `WebFilter`
- Mauricio ECR
- Snippets
- 30 Aug, 2025
En entornos modernos de microservicios, la consistencia en las respuestas de una API es más que una cuestión de estética: es un factor crítico para la mantenibilidad, observabilidad y experiencia del consumidor. Aplicaciones frontend, integraciones con terceros, herramientas de monitoreo y otros microservicios esperan estructuras de respuesta predecibles. Cada variación no planificada introduce fricción: más lógica en los clientes, validaciones dispersas y puntos ciegos en trazabilidad.
En este contexto, estandarizar las respuestas de manera transversal —sin ensuciar cada controlador con lógica repetitiva— no solo simplifica el desarrollo, también abre la puerta a métricas uniformes, trazabilidad distribuida y soporte para nuevas funcionalidades sin tocar el código de negocio.
Este artículo explica cómo lograrlo en aplicaciones reactivas con Spring WebFlux, donde la naturaleza streaming de la respuesta introduce desafíos distintos a los de un stack imperativo como Spring MVC.
El Contrato de Respuesta: Mucho más que Datos
Antes de modificar nada, debemos definir el destino. Una respuesta estándar debe separar claramente los datos de negocio de la información contextual que permite entender la petición en su conjunto.
Un diseño común y extensible puede lucir así:
package com.app247.api.shared.response_wrapper.model;
import lombok.Builder;
import lombok.Data;
@Data
@Builder
public class ApiResponse<T> {
private Meta meta;
private T data;
@Data
@Builder
public static class Meta {
private String timestamp;
private String path;
private int status;
private String requestId;
}
}
Este contrato permite:
- Consistencia: cada respuesta, sin importar el endpoint, sigue la misma forma.
- Trazabilidad: con
requestIdytimestamppodemos correlacionar logs, métricas y reportes. - Extensibilidad: podemos agregar campos en
meta(e.g., tiempos de respuesta, versión del servicio) sin afectar al cliente.
En entornos con OpenAPI/Swagger, este modelo puede documentarse fácilmente para que los consumidores conozcan el formato exacto de las respuestas.
WebFlux y el Desafío del Streaming
En aplicaciones no reactivas, ResponseBodyAdvice permite interceptar y modificar respuestas antes de serializarse. Pero en WebFlux, las respuestas son streams (Publisher<DataBuffer>), no objetos finales en memoria.
Esto implica dos retos:
- Respetar el modelo reactivo: no bloquear el flujo ni forzar materializaciones tempranas.
- Actuar en el punto correcto: cuando la respuesta está completa, pero antes de enviarla al cliente.
Aquí entra en juego el dúo WebFilter + ServerHttpResponseDecorator. El filtro decide si aplicar la transformación; el decorador define cómo hacerlo.
El Filtro: Decidiendo Cuándo Intervenir
Nuestro WebFilter actúa como middleware, excluyendo rutas (por ejemplo, Swagger o Actuator) y habilitando/deshabilitando la lógica según configuración externa:
package com.app247.api.shared.response_wrapper.filter;
import com.app247.api.shared.response_wrapper.config.ResponseWrapperProperties;
import com.app247.api.shared.response_wrapper.decorator.ResponseWrapperDecorator;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.core.annotation.Order;
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
import org.springframework.stereotype.Component;
import org.springframework.util.AntPathMatcher;
import org.springframework.web.server.ServerWebExchange;
import org.springframework.web.server.WebFilter;
import org.springframework.web.server.WebFilterChain;
import reactor.core.publisher.Mono;
@Slf4j
@Component
@Order(-2)
@RequiredArgsConstructor
public class ResponseWrapperFilter implements WebFilter {
private final ResponseWrapperProperties properties;
private final ObjectMapper objectMapper;
private final AntPathMatcher pathMatcher = new AntPathMatcher();
@Override
public Mono<Void> filter(ServerWebExchange exchange, WebFilterChain chain) {
String path = exchange.getRequest().getURI().getPath();
// Verificamos si la ruta está excluida
boolean isExcluded = !properties.isEnabled() || properties.getExcludedPaths().stream()
.anyMatch(pattern -> pathMatcher.match(pattern, path));
if (isExcluded) {
return chain.filter(exchange);
}
// Creamos una instancia de nuestro nuevo decorador
ServerHttpResponseDecorator decoratedResponse = new ResponseWrapperDecorator(
exchange.getResponse(),
path,
objectMapper
);
// Pasamos el exchange con la respuesta decorada al siguiente filtro en la cadena
return chain.filter(exchange.mutate().response(decoratedResponse).build());
}
}
Las rutas excluidas y la activación del wrapper se controlan con propiedades externas, evitando recompilar para cambios operativos:
package com.app247.api.shared.response_wrapper.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
import java.util.ArrayList;
import java.util.List;
/*
Ejemplo:
api:
response:
wrapper:
enabled: true
# Patrones de URL para excluir. Usa el formato Ant.
excluded-paths:
- "/v3/api-docs/**"
- "/swagger-ui/**"
- "/webjars/**"
- "/swagger-resources/**"
- "/actuator/**"
*/
@Data
@Component
@ConfigurationProperties(prefix = "api.response.wrapper")
public class ResponseWrapperProperties {
private boolean enabled = true;
private List<String> excludedPaths = new ArrayList<>();
}
El Decorador: Interviniendo sin Romper el Flujo
ServerHttpResponseDecorator nos da acceso al cuerpo de la respuesta. El método clave es writeWith, que recibe el stream de datos antes de enviarlo al cliente.
package com.app247.api.shared.response_wrapper.decorator;
import com.app247.api.shared.response_wrapper.model.ApiResponse;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.reactivestreams.Publisher;
import org.springframework.core.io.buffer.DataBuffer;
import org.springframework.core.io.buffer.DataBufferUtils;
import org.springframework.core.io.buffer.DefaultDataBufferFactory;
import org.springframework.http.HttpStatus;
import org.springframework.http.HttpStatusCode;
import org.springframework.http.MediaType;
import org.springframework.http.server.reactive.ServerHttpResponse;
import org.springframework.http.server.reactive.ServerHttpResponseDecorator;
import reactor.core.publisher.Mono;
import java.nio.charset.StandardCharsets;
import java.time.Instant;
import java.util.UUID;
/**
* Decorador para ServerHttpResponse que intercepta las respuestas exitosas
* y las envuelve en una estructura estandarizada de ApiResponse (meta y data).
*/
@Slf4j
public class ResponseWrapperDecorator extends ServerHttpResponseDecorator {
private final ObjectMapper objectMapper;
private final String path;
public ResponseWrapperDecorator(ServerHttpResponse delegate, String path, ObjectMapper objectMapper) {
super(delegate);
this.path = path;
this.objectMapper = objectMapper;
}
/**
* Sobrescribe el método que escribe el cuerpo de la respuesta en el flujo de salida.
* Aquí es donde ocurre toda la magia de la intercepción y transformación.
* @param body El publicador original del cuerpo de la respuesta.
* @return Un Mono<Void> que representa la finalización de la operación de escritura.
*/
@Override
public Mono<Void> writeWith(Publisher<? extends DataBuffer> body) {
// PASO 1: Almacenar el cuerpo completo en un búfer.
// DataBufferUtils.join() consume tod_o el flujo del 'body' y lo une en un solo DataBuffer.
// Esto es CRUCIAL porque crea un punto de sincronización. La lógica siguiente
// no se ejecutará hasta que el controlador haya terminado y el cuerpo completo esté disponible.
Mono<DataBuffer> bufferedBody = DataBufferUtils.join(body)
.defaultIfEmpty(new DefaultDataBufferFactory().wrap(new byte[0])); // Maneja cuerpos vacíos (ej: 204 No Content)
// PASO 2: Usar flatMap para transformar el cuerpo almacenado en búfer.
// El código dentro de flatMap está garantizado a ejecutarse DESPUÉS de que 'bufferedBody' se complete.
return bufferedBody.flatMap(originalBuffer -> {
// PASO 3: Obtener el código de estado.
// En este punto, la llamada a getStatusCode() es 100% fiable porque el controlador
// ya ha finalizado y el framework ha establecido el estado final de la respuesta.
HttpStatusCode statusCode = getStatusCode();
// PASO 4: Decidir si se debe envolver la respuesta.
// Si el estado es un error explícito (4xx o 5xx), no hacemos nada y devolvemos el cuerpo original.
if (statusCode != null && !statusCode.is2xxSuccessful()) {
// Se escribe el buffer original en la respuesta real.
return getDelegate().writeWith(Mono.just(originalBuffer));
}
// PASO 5: Manejar el caso del entorno de pruebas.
// En WebFluxTest, un 200 OK por defecto puede resultar en un statusCode 'null'.
// Asumimos HttpStatus.OK si el estado es null para que las pruebas pasen.
HttpStatusCode statusToUse = (statusCode != null) ? statusCode : HttpStatus.OK;
// PASO 6: Procesar y envolver el cuerpo de la respuesta.
byte[] bytes = new byte[originalBuffer.readableByteCount()];
originalBuffer.read(bytes);
DataBufferUtils.release(originalBuffer); // Liberar memoria del buffer original.
String originalBodyJson = new String(bytes, StandardCharsets.UTF_8);
// Evitar envolver una respuesta que ya tiene nuestro formato.
if (originalBodyJson.contains("\"meta\"")) {
return getDelegate().writeWith(Mono.just(new DefaultDataBufferFactory().wrap(bytes)));
}
try {
// Deserializar el cuerpo original para poder ponerlo dentro del campo 'data'.
// Si el cuerpo está vacío, se asigna 'null' a los datos.
Object originalBodyObject = originalBodyJson.isEmpty() ? null : objectMapper.readValue(originalBodyJson, Object.class);
// Construir la nueva respuesta envuelta.
ApiResponse<?> apiResponse = buildSuccessResponse(originalBodyObject, path, statusToUse);
// Serializar la respuesta envuelta a bytes.
byte[] responseBytes = objectMapper.writeValueAsBytes(apiResponse);
// Actualizar las cabeceras HTTP con la nueva longitud y tipo de contenido.
getHeaders().setContentLength(responseBytes.length);
getHeaders().setContentType(MediaType.APPLICATION_JSON);
// Crear un nuevo buffer con la respuesta envuelta.
DataBuffer wrappedBuffer = new DefaultDataBufferFactory().wrap(responseBytes);
// Escribir el nuevo cuerpo en la respuesta real. Esta es la llamada final y única
// que envía los datos al cliente, siguiendo las buenas prácticas reactivas.
return getDelegate().writeWith(Mono.just(wrappedBuffer));
} catch (Exception e) {
log.error("Error al envolver la respuesta para la ruta {}: {}", path, e.getMessage(), e);
return getDelegate().writeWith(Mono.just(new DefaultDataBufferFactory().wrap(bytes)));
}
});
}
/**
* Método de ayuda para construir la estructura estandarizada de ApiResponse.
* @param data El objeto de datos original que se incluirá en el campo 'data'.
* @param path La ruta de la petición actual.
* @param status El código de estado HTTP final.
* @return Una instancia de ApiResponse.
*/
private ApiResponse<?> buildSuccessResponse(Object data, String path, HttpStatusCode status) {
ApiResponse.Meta meta = ApiResponse.Meta.builder()
.timestamp(Instant.now().toString())
.path(path)
.requestId(UUID.randomUUID().toString().substring(0, 10))
.status(status.value())
.build();
return ApiResponse.builder()
.meta(meta)
.data(data)
.build();
}
}
Consideraciones Técnicas
- Performance:
DataBufferUtils.join()carga todo en memoria; para respuestas muy grandes, conviene evaluar streaming JSON. - Idempotencia: el filtro detecta si ya existe
"meta"para evitar doble envoltura. - Trazabilidad distribuida:
requestIdpuede integrarse con Spring Cloud Sleuth o MDC para correlacionar logs entre microservicios.
Pruebas: Validando Comportamiento y Robustez
Con @WebFluxTest podemos probar controladores y filtros en un entorno aislado.
package com.app247.api.shared.response_wrapper.filter;
import com.app247.api.shared.response_wrapper.config.ResponseWrapperProperties;
import com.app247.api.shared.response_wrapper.model.ApiResponse;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.autoconfigure.web.reactive.WebFluxTest;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.springframework.test.web.reactive.server.WebTestClient;
import org.springframework.util.AntPathMatcher;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.time.Instant;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import static org.mockito.Mockito.when;
/**
* Pruebas de integración para el ResponseWrapperFilter.
* Valida que las respuestas exitosas se envuelvan y que las de error o excluidas se ignoren.
*/
@WebFluxTest
@Import(ResponseWrapperFilterTest.TestConfig.class)
class ResponseWrapperFilterTest {
@Autowired
private WebTestClient webTestClient;
@MockitoBean
private ResponseWrapperProperties responseWrapperProperties;
@BeforeEach
void setUp() {
when(responseWrapperProperties.isEnabled()).thenReturn(true);
when(responseWrapperProperties.getExcludedPaths()).thenReturn(List.of("/excluded/**"));
}
@Test
@DisplayName("Debería envolver una respuesta Mono exitosa en ApiResponse")
void shouldWrapSuccessfulMonoResponse() {
webTestClient.get().uri("/test/mono")
.exchange()
.expectStatus().isOk()
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
.jsonPath("$.meta").exists()
.jsonPath("$.meta.status").isEqualTo(200)
.jsonPath("$.meta.path").isEqualTo("/test/mono")
.jsonPath("$.data").exists()
.jsonPath("$.data.id").isEqualTo(1)
.jsonPath("$.data.name").isEqualTo("Test Mono");
}
@Test
@DisplayName("Debería envolver una respuesta Flux exitosa en ApiResponse con una lista")
void shouldWrapSuccessfulFluxResponse() {
webTestClient.get().uri("/test/flux")
.exchange()
.expectStatus().isOk()
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
.jsonPath("$.meta").exists()
.jsonPath("$.meta.status").isEqualTo(200)
.jsonPath("$.data").isArray()
.jsonPath("$.data[0].id").isEqualTo(1)
.jsonPath("$.data[0].name").isEqualTo("Test Flux 1")
.jsonPath("$.data[1].id").isEqualTo(2)
.jsonPath("$.data[1].name").isEqualTo("Test Flux 2");
}
@Test
@DisplayName("No debería envolver una respuesta de una ruta excluida")
void shouldNotWrapExcludedPath() {
webTestClient.get().uri("/excluded/path")
.exchange()
.expectStatus().isOk()
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
.jsonPath("$.meta").doesNotExist()
.jsonPath("$.data").doesNotExist()
.jsonPath("$.id").isEqualTo(99)
.jsonPath("$.name").isEqualTo("Excluded");
}
@Test
@DisplayName("No debería envolver una respuesta de error (ej: 400 Bad Request)")
void shouldNotWrapErrorResponse() {
webTestClient.get().uri("/test/error")
.exchange()
.expectStatus().isBadRequest()
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
.jsonPath("$.meta").exists()
.jsonPath("$.errors").exists()
.jsonPath("$.errors[0].code").isEqualTo("400-CUSTOM-ERROR")
.jsonPath("$.data").doesNotExist();
}
@Test
@DisplayName("No debería envolver una respuesta que ya tiene el formato ApiResponse")
void shouldNotDoubleWrapAlreadyFormattedResponse() {
webTestClient.get().uri("/test/pre-wrapped")
.exchange()
.expectStatus().isOk()
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
.jsonPath("$.meta").exists()
.jsonPath("$.meta.status").isEqualTo(200)
.jsonPath("$.data.message").isEqualTo("This is already wrapped")
.jsonPath("$.data.meta").doesNotExist(); // La comprobación clave: no hay un 'meta' dentro del 'data'.
}
@Test
@DisplayName("Debería devolver un error 500 estándar si la serialización del framework falla")
void shouldReturnStandard500ErrorOnFrameworkSerializationFailure() {
webTestClient.get().uri("/test/unserializable")
.exchange()
// 1. Aserción clave: el estado DEBE ser 500 Internal Server Error.
.expectHeader().contentType(MediaType.APPLICATION_JSON)
.expectBody()
// 2. Aserciones sobre el cuerpo de error estándar de Spring Boot.
// Este cuerpo NO es el original, sino el generado por el manejador de errores de Spring.
.jsonPath("$.status").isEqualTo(500)
.jsonPath("$.error").isEqualTo("Internal Server Error")
.jsonPath("$.path").isEqualTo("/test/unserializable")
// 3. Confirmamos que no hay rastro del cuerpo original ni de nuestra envoltura personalizada.
.jsonPath("$.meta").doesNotExist()
.jsonPath("$.data").doesNotExist()
.jsonPath("$.id").doesNotExist();
}
// --- CONFIGURACIÓN INTERNA Y COMPONENTES DE PRUEBA ---
@Data
@NoArgsConstructor
@AllArgsConstructor
static class TestDto {
private int id;
private String name;
}
// DTO diseñado para fallar durante la serialización de Jackson debido a una referencia circular.
@Data
static class UnserializableDto {
private int id = 123;
private Object problematicField = this;
}
@RestController
static class TestController {
@GetMapping("/test/mono")
Mono<TestDto> getMono() {
return Mono.just(new TestDto(1, "Test Mono"));
}
@GetMapping("/test/flux")
Flux<TestDto> getFlux() {
return Flux.just(new TestDto(1, "Test Flux 1"), new TestDto(2, "Test Flux 2"));
}
@GetMapping("/excluded/path")
Mono<TestDto> getExcluded() {
return Mono.just(new TestDto(99, "Excluded"));
}
@GetMapping("/test/error")
Mono<TestDto> getError() {
return Mono.error(new BusinessException("Error forzado", "400-CUSTOM-ERROR"));
}
@GetMapping("/test/pre-wrapped")
Mono<ApiResponse<Map<String, String>>> getPreWrappedResponse() {
ApiResponse.Meta meta = ApiResponse.Meta.builder().status(200).build();
Map<String, String> data = Collections.singletonMap("message", "This is already wrapped");
return Mono.just(ApiResponse.<Map<String, String>>builder().meta(meta).data(data).build());
}
@GetMapping("/test/unserializable")
Mono<UnserializableDto> getUnserializableObject() {
return Mono.just(new UnserializableDto());
}
}
static class BusinessException extends RuntimeException {
private final String errorCode;
public BusinessException(String message, String errorCode) {
super(message);
this.errorCode = errorCode;
}
public String getErrorCode() {
return errorCode;
}
}
@org.springframework.web.bind.annotation.RestControllerAdvice
static class TestGlobalExceptionHandler {
@org.springframework.web.bind.annotation.ExceptionHandler(BusinessException.class)
@org.springframework.web.bind.annotation.ResponseStatus(HttpStatus.BAD_REQUEST)
public Mono<Map<String, Object>> handleBusinessException(BusinessException ex) {
Map<String, String> error = Map.of("code", ex.getErrorCode(), "message", ex.getMessage());
Map<String, Object> meta = Map.of("timestamp", Instant.now().toString());
return Mono.just(Map.of("meta", meta, "errors", List.of(error)));
}
}
@Configuration
static class TestConfig {
@Bean
public ObjectMapper objectMapper() {
return new ObjectMapper();
}
@Bean
public AntPathMatcher antPathMatcher() {
return new AntPathMatcher();
}
@Bean
public TestGlobalExceptionHandler testGlobalExceptionHandler() {
return new TestGlobalExceptionHandler();
}
@Bean
public ResponseWrapperFilter responseWrapperFilter(
ResponseWrapperProperties properties, ObjectMapper objectMapper
) {
return new ResponseWrapperFilter(properties, objectMapper);
}
@Bean
public TestController testController() {
return new TestController();
}
}
}
Podemos extender las pruebas con StepVerifier para validar que el flujo sigue siendo reactivo y no introduce bloqueos inesperados.
Próximos Pasos y Extensiones
La solución presentada puede evolucionar hacia:
- Trazabilidad distribuida: Propagando
requestIdcon Spring Cloud Sleuth, Zipkin o Jaeger. - Internacionalización: Soporte para mensajes localizados en errores o advertencias.
- Observabilidad avanzada: Tiempo de procesamiento en
meta, integración con Prometheus o Grafana. - Functional Endpoints: Adaptando la solución a APIs basadas en
RouterFunctionen lugar de anotaciones tradicionales.
Con esta base, la envoltura de respuestas deja de ser solo un detalle de formato y se convierte en una capa estratégica para consistencia, trazabilidad y mantenimiento a largo plazo.
Transacciones Declarativas en Arquitecturas Hexagonales con Spring WebFlux y AOP
- Mauricio ECR
- Snippets
- 27 Aug, 2025
En el desarrollo de aplicaciones reactivas modernas, especialmente bajo paradigmas como la Arquitectura Hexagonal, surgen desafíos que nos obligan a repensar cómo aplicamos conceptos transversales. Un
Transacciones Declarativas en Arquitecturas Hexagonales con Spring WebFlux y AOP
- Mauricio ECR
- Snippets
- 27 Aug, 2025
En el desarrollo de aplicaciones reactivas modernas, especialmente bajo paradigmas como la Arquitectura Hexagonal, surgen desafíos que nos obligan a repensar cómo aplicamos conceptos transversales. Uno de los interrogantes más comunes es: ¿dónde y cómo gestionamos las transacciones de base de datos sin contaminar nuestra lógica de negocio? Este artículo documenta un viaje desde esa pregunta inicial hasta una solución robusta y elegante, utilizando el poder de la Programación Orientada a Aspectos (AOP) en un entorno Spring WebFlux con R2DBC.
La Arquitectura como Punto de Partida
Antes de sumergirnos en el código, es fundamental visualizar la estructura del proyecto. Una organización clara de paquetes, que refleje las capas de la Arquitectura Hexagonal, es la base sobre la que construiremos nuestra solución. El dominio permanece en el centro, puro y sin dependencias externas, mientras que la aplicación y la infraestructura se organizan a su alrededor.
ms_auth/
├── applications/app-service/ # Módulo principal de la aplicación Spring Boot
│ ├── build.gradle
│ └── src/
│ ├── main/java/com/app247/
│ │ ├── MainApplication.java
│ │ └── config/aop/
│ │ └── TransactionalUseCaseAspect.java # Nuestro Aspecto AOP
│ └── test/java/com/app247/config/aop/
│ ├── TransactionalUseCaseAspectTest.java # Test unitario del Aspecto
│ └── TransactionalRollbackSelfContainedTest.java # Test de Integración
│
├── domain/
│ ├── model/
│ └── usecase/ # Módulo de la lógica de negocio pura
│ └── src/main/java/com/app247/usecase/shared/core/usecase/
│ ├── TransactionalWrapperUseCase.java # Anotación personalizada
│ └── UseCase.java # Interfaz genérica
│
└── infrastructure/
├── r2dbc-postgresql/ # Módulo adaptador para la base de datos
└── reactive-web/ # Módulo adaptador para los controladores REST
El Dilema Inicial: La Transacción y la Unidad de Trabajo
Todo comienza con una necesidad fundamental: asegurar la atomicidad de las operaciones. Imaginemos un caso de uso de negocio, como procesar una compra, que implica modificar el inventario de productos y crear un registro de orden. Ambas acciones deben tener éxito, o ninguna debe persistir. Esta es la definición de una unidad de trabajo, y la herramienta para garantizarla es la transacción.
La primera intuición podría ser colocar la anotación @Transactional de Spring en los métodos del repositorio. Sin embargo, esto es incorrecto. Una transacción en el repositorio solo cubriría una única operación de base de datos, rompiendo la unidad de trabajo del negocio. La transacción debe envolver la ejecución completa del caso de uso.
Esto nos lleva a la capa de servicio o caso de uso. Pero aquí nos encontramos con el primer gran obstáculo arquitectónico. En una Arquitectura Hexagonal, la capa de dominio (donde residen los casos de uso) debe ser pura. No puede, ni debe, tener dependencias de frameworks externos como Spring. Anotar un caso de uso del dominio con @Transactional viola este principio fundamental, acoplando nuestra lógica de negocio más preciada a un detalle de infraestructura.
La Solución Emerge: Programación Orientada a Aspectos
Si no podemos modificar el dominio, debemos aplicar el comportamiento transaccional desde afuera, de una manera no invasiva. Aquí es donde la Programación Orientada a Aspectos (AOP) brilla. AOP nos permite interceptar la ejecución de nuestros métodos para añadir funcionalidades transversales (como transacciones, seguridad o logging) sin alterar el código original.
La estrategia que emerge es crear un mecanismo declarativo y reutilizable que nos permita "marcar" qué casos de uso deben ser transaccionales, dejando que la magia de AOP haga el resto.
Una Anotación para Declarar la Intención
El primer paso es crear una anotación personalizada. Su único propósito es servir como una señal o marcador. Al ser parte de nuestro código de dominio (usecase), no introduce una dependencia directa de Spring, sino que define un contrato interno.
package com.app247.usecase.shared.core.usecase;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
* Anotación para marcar clases de Casos de Uso que deben ser
* envueltas en una transacción reactiva de forma automática.
*/
@Target(ElementType.TYPE) // Se aplica a nivel de clase
@Retention(RetentionPolicy.RUNTIME) // Disponible en tiempo de ejecución para que Spring la lea
public @interface TransactionalWrapperUseCase {
}
Junto a esta, podemos definir una interfaz genérica para estandarizar nuestros casos de uso, promoviendo un diseño limpio y consistente.
package com.app247.usecase.shared.core.usecase;
// Interfaz genérica (opcional pero recomendada)
public interface UseCase<Request, Response> {
Response execute(Request request);
}
El Aspecto: El Motor de la Transacción
Con la anotación en su lugar, construimos el componente que buscará esta marca y aplicará la lógica transaccional. Este es nuestro Aspecto, una clase de infraestructura que vive en la capa de aplicación.
Este Aspecto tiene dos partes clave:
- Pointcut: Una expresión que actúa como un selector. Le dice a Spring: "Encuentra todos los métodos públicos en cualquier clase que esté anotada con
@TransactionalWrapperUseCase". - Advice: La lógica que se ejecuta cuando el Pointcut encuentra una coincidencia. Usaremos un
advicede tipo@Around, que nos permite envolver completamente la ejecución del método original.
La lógica del advice es simple pero poderosa: toma el Mono o Flux devuelto por el caso de uso y lo compone con el TransactionalOperator reactivo de Spring. Este operador se encarga de iniciar la transacción antes de la suscripción y de realizar commit o rollback al finalizar.
package com.app247.config.aop;
import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.annotation.Pointcut;
import org.springframework.stereotype.Component;
import org.springframework.transaction.reactive.TransactionalOperator;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@Aspect
@Component
public class TransactionalUseCaseAspect {
private final TransactionalOperator transactionalOperator;
public TransactionalUseCaseAspect(TransactionalOperator transactionalOperator) {
this.transactionalOperator = transactionalOperator;
}
@Pointcut("@within(com.app247.usecase.shared.core.usecase.TransactionalWrapperUseCase) && execution(public * *(..))")
public void transactionalUseCase() {
// Método vacío para nombrar el pointcut.
}
@Around("transactionalUseCase()")
public Object wrapInTransaction(ProceedingJoinPoint joinPoint) throws Throwable {
Object result = joinPoint.proceed();
if (result instanceof Mono) {
return ((Mono<?>) result).as(transactionalOperator::transactional);
} else if (result instanceof Flux) {
return ((Flux<?>) result).as(transactionalOperator::transactional);
}
return result;
}
}
Con estos dos elementos, hemos creado un sistema donde simplemente anotando una clase de caso de uso con @TransactionalWrapperUseCase, garantizamos que su ejecución será atómica, sin haber escrito una sola línea de código transaccional dentro del propio caso de uso.
Probando la Solución: De la Confianza a la Certeza
Una solución no está completa hasta que se prueba rigurosamente. Para este mecanismo, necesitamos dos niveles de prueba para tener una confianza total.
Nivel 1: El Test de Cableado (Unitario)
El primer test debe responder a la pregunta: ¿Nuestro aspecto AOP está correctamente configurado para interceptar la llamada y usar el TransactionalOperator? Este test valida tanto respuestas Mono como Flux.
Este test no necesita una base de datos. Utiliza un contexto de Spring para activar el mecanismo AOP, pero reemplaza todas las dependencias externas (TransactionalOperator, repositorios) con Mocks. El objetivo no es probar el rollback, sino verificar la interacción: que el método transactional() del operador sea invocado.
package com.app247.config.aop;
import com.app247.usecase.shared.core.usecase.TransactionalWrapperUseCase;
import org.junit.jupiter.api.Test;
import org.mockito.InjectMocks;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.EnableAspectJAutoProxy;
import org.springframework.context.annotation.Import;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.springframework.transaction.reactive.TransactionalOperator;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import static org.assertj.core.api.AssertionsForClassTypes.assertThat;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.*;
@SpringBootTest(classes = TransactionalUseCaseAspectTest.TestConfig.class)
class TransactionalUseCaseAspectTest {
@Autowired
private PurchaseProductUseCasePort purchaseUseCase;
@Autowired
private FindProductsUseCasePort findProductsUseCase; // Caso de uso que devuelve Flux
@Autowired
private NonReactiveUseCasePort nonReactiveUseCase;
@MockitoBean
private ProductRepository productRepository;
@MockitoBean
private OrderRepository orderRepository;
@MockitoBean
private TransactionalOperator transactionalOperator;
@InjectMocks
private TransactionalUseCaseAspect transactionalUseCaseAspect;
@Test
void whenUseCaseReturnsMono_thenItShouldBeWrappedInTransaction() {
// ARRANGE
Product fakeProduct = new Product("prod-123", 10);
Order fakeOrder = new Order("user-007", "prod-123");
when(productRepository.findById(any())).thenReturn(Mono.just(fakeProduct));
when(orderRepository.save(any())).thenReturn(Mono.just(fakeOrder));
when(productRepository.updateStock(any(), any(Integer.class))).thenReturn(Mono.empty());
when(transactionalOperator.transactional(any(Mono.class)))
.thenAnswer(invocation -> invocation.getArgument(0));
// ACT
Mono<Order> result = purchaseUseCase.execute("user-007", "prod-123");
// ASSERT
StepVerifier.create(result).expectNext(fakeOrder).verifyComplete();
verify(transactionalOperator).transactional(any(Mono.class));
verify(productRepository).updateStock("prod-123", 9);
}
@Test
void whenUseCaseReturnsFlux_thenItShouldBeWrappedInTransaction() {
// ARRANGE
Product fakeProduct1 = new Product("prod-001", 5);
Product fakeProduct2 = new Product("prod-002", 3);
when(productRepository.findAll()).thenReturn(Flux.just(fakeProduct1, fakeProduct2));
// Configuramos el mock para que el operador transaccional simplemente devuelva el Flux original
when(transactionalOperator.transactional(any(Flux.class)))
.thenAnswer(invocation -> invocation.getArgument(0));
// ACT
Flux<Product> result = findProductsUseCase.execute(null); // `null` porque no requiere parámetros
// ASSERT
StepVerifier.create(result)
.expectNext(fakeProduct1)
.expectNext(fakeProduct2)
.verifyComplete();
// La verificación clave: ¿Se llamó al operador con un Flux?
verify(transactionalOperator).transactional(any(Flux.class));
}
/**
* Test para el caso no reactivo.
*/
@Test
void whenUseCaseIsNotReactive_thenItShouldNotBeWrappedInTransaction() {
// --- ARRANGE (Preparar) ---
String expectedResult = "Este es un resultado síncrono";
// --- ACT (Actuar) ---
// Ejecutamos el caso de uso que devuelve un String simple.
String actualResult = nonReactiveUseCase.execute(null);
// --- ASSERT (Verificar) ---
// 1. Verificamos que el resultado devuelto es el original, sin cambios.
assertThat(actualResult).isEqualTo(expectedResult);
// 2. La verificación MÁS IMPORTANTE: nos aseguramos de que el operador transaccional
// NUNCA fue invocado, ya que la respuesta no era ni Mono ni Flux.
verify(transactionalOperator, never()).transactional(any(Mono.class));
verify(transactionalOperator, never()).transactional(any(Flux.class));
}
@Test
void transactionalUseCasePointcut_shouldExecuteForCoverage() {
// --- ACT ---
// Simplemente llamamos al método vacío.
// La herramienta de cobertura registrará que se ha entrado en este método.
// --- ASSERT ---
// Como el método no hace nada, la única aserción posible es
// que la llamada no lance ninguna excepción.
assertDoesNotThrow(() -> {
transactionalUseCaseAspect.transactionalUseCase();
});
}
@Configuration
@EnableAspectJAutoProxy
@Import(TransactionalUseCaseAspect.class)
static class TestConfig {
@Bean
public PurchaseProductUseCasePort purchaseProductUseCase(ProductRepository productRepo, OrderRepository orderRepo) {
return new PurchaseProductUseCase(productRepo, orderRepo);
}
@Bean
public FindProductsUseCasePort findProductsUseCase(ProductRepository productRepo) {
return new FindProductsUseCase(productRepo);
}
@Bean
public NonReactiveUseCasePort nonReactiveUseCase() {
return new NonReactiveUseCase();
}
}
// --- Definiciones Fakes ---
record Product(String id, int stock) {}
record Order(String userId, String productId) {}
interface ProductRepository {
Mono<Product> findById(String productId);
Flux<Product> findAll(); // Añadido para el test de Flux
Mono<Void> updateStock(String productId, int newStock);
}
interface OrderRepository { Mono<Order> save(Order order); }
interface PurchaseProductUseCasePort { Mono<Order> execute(String userId, String productId); }
interface FindProductsUseCasePort { Flux<Product> execute(Void request); } // Nuevo caso de uso para Flux
interface NonReactiveUseCasePort { String execute(Void request); }
@TransactionalWrapperUseCase
static class PurchaseProductUseCase implements PurchaseProductUseCasePort {
private final ProductRepository productRepository;
private final OrderRepository orderRepository;
public PurchaseProductUseCase(ProductRepository p, OrderRepository o) { this.productRepository = p; this.orderRepository = o; }
public Mono<Order> execute(String userId, String productId) {
return productRepository.findById(productId)
.flatMap(product -> productRepository.updateStock(product.id(), product.stock() - 1)
.then(orderRepository.save(new Order(userId, productId))));
}
}
@TransactionalWrapperUseCase
static class FindProductsUseCase implements FindProductsUseCasePort {
private final ProductRepository productRepository;
public FindProductsUseCase(ProductRepository p) { this.productRepository = p; }
public Flux<Product> execute(Void request) {
return productRepository.findAll();
}
}
@TransactionalWrapperUseCase
static class NonReactiveUseCase implements NonReactiveUseCasePort {
@Override
public String execute(Void request) {
return "Este es un resultado síncrono";
}
}
}
Nivel 2: El Test de Comportamiento (Integración)
El segundo test debe responder a una pregunta más importante: si una operación falla, ¿la transacción realmente hace rollback?
Para esto, necesitamos un test de integración que utilice una base de datos real (en memoria, como H2, para velocidad y aislamiento) y el TransactionalOperator real de Spring. La clave aquí es usar @SpyBean para envolver un repositorio real y forzar un fallo en una de sus operaciones. La validación final consiste en consultar la base de datos después del fallo y verificar que el estado de los datos ha sido revertido a su estado original.
Este test es completamente autocontenido: define su propia configuración, su esquema de base de datos y sus implementaciones de dominio e infraestructura, pero lo más importante es que importa y prueba el Aspecto de AOP de producción real.
package com.app247.config.aop;
import com.app247.usecase.shared.core.usecase.TransactionalWrapperUseCase;
import io.r2dbc.spi.ConnectionFactory;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.boot.autoconfigure.context.PropertyPlaceholderAutoConfiguration;
import org.springframework.boot.autoconfigure.r2dbc.R2dbcAutoConfiguration;
import org.springframework.boot.autoconfigure.transaction.TransactionAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.EnableAspectJAutoProxy;
import org.springframework.context.annotation.Import;
import org.springframework.data.annotation.Id;
import org.springframework.data.r2dbc.core.R2dbcEntityTemplate;
import org.springframework.data.relational.core.mapping.Table;
import org.springframework.r2dbc.connection.R2dbcTransactionManager;
import org.springframework.r2dbc.core.DatabaseClient;
import org.springframework.stereotype.Repository;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.TestPropertySource;
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.doReturn;
import static org.springframework.data.relational.core.query.Criteria.where;
import static org.springframework.data.relational.core.query.Query.query;
@SpringBootTest(classes = TransactionalRollbackSelfContainedTest.TestConfig.class)
@ImportAutoConfiguration({
R2dbcAutoConfiguration.class,
TransactionAutoConfiguration.class,
PropertyPlaceholderAutoConfiguration.class
})
@TestPropertySource(properties = {
"spring.r2dbc.url=r2dbc:h2:mem:///finaltestdb;DB_CLOSE_DELAY=-1;",
"spring.r2dbc.username=sa",
"spring.r2dbc.password=",
"spring.sql.init.mode=never"
})
@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
@TestInstance(TestInstance.Lifecycle.PER_CLASS)
class TransactionalRollbackSelfContainedTest {
@Autowired private PurchaseProductUseCasePort purchaseUseCase;
@Autowired private DatabaseClient databaseClient;
@Autowired private R2dbcEntityTemplate template;
@MockitoSpyBean
private OrderRepository orderRepository;
private final String PRODUCT_ID = "prod-123";
private final int INITIAL_STOCK = 10;
@BeforeAll
void setupDatabaseSchema() {
String createProductsTable = "CREATE TABLE PRODUCTS (id VARCHAR(255) PRIMARY KEY, name VARCHAR(255), stock INT);";
String createOrdersTable = "CREATE TABLE ORDERS (id INT AUTO_INCREMENT PRIMARY KEY, user_id VARCHAR(255), product_id VARCHAR(255));";
databaseClient.sql(createProductsTable).then().block();
databaseClient.sql(createOrdersTable).then().block();
}
@BeforeEach
void setupTestData() {
databaseClient.sql("DELETE FROM PRODUCTS").then().block();
template.insert(new ProductEntity(PRODUCT_ID, "Test Product", INITIAL_STOCK)).block();
}
@Test
void whenSecondOperationFails_thenRealAspectRollsBackTransaction() {
doReturn(Mono.error(new RuntimeException("DB Error"))).when(orderRepository).save(any());
Mono<Void> result = purchaseUseCase.execute("user-007", PRODUCT_ID);
StepVerifier.create(result).expectError(RuntimeException.class).verify();
ProductEntity productAfter = template.selectOne(query(where("id").is(PRODUCT_ID)), ProductEntity.class).block();
assertThat(productAfter.stock()).isEqualTo(INITIAL_STOCK);
}
@Configuration
@EnableAspectJAutoProxy
@Import(TransactionalUseCaseAspect.class)
static class TestConfig {
@Bean
public R2dbcEntityTemplate r2dbcEntityTemplate(ConnectionFactory connectionFactory) {
return new R2dbcEntityTemplate(connectionFactory);
}
@Bean
public R2dbcTransactionManager transactionManager(ConnectionFactory connectionFactory) {
return new R2dbcTransactionManager(connectionFactory);
}
@Bean
public PurchaseProductUseCasePort purchaseProductUseCase(ProductRepository productRepo, OrderRepository orderRepo) {
return new PurchaseProductUseCase(productRepo, orderRepo);
}
@Bean
public ProductRepository productRepository(R2dbcEntityTemplate template) {
return new R2dbcProductRepositoryAdapter(template);
}
@Bean
public OrderRepository orderRepository(R2dbcEntityTemplate template) {
return new R2dbcOrderRepositoryAdapter(template);
}
}
record Product(String id, int stock) {}
record Order(String userId, String productId) {}
interface ProductRepository { Mono<Product> findById(String id); Mono<Void> updateStock(String id, int stock); }
interface OrderRepository { Mono<Order> save(Order order); }
interface PurchaseProductUseCasePort { Mono<Void> execute(String userId, String productId); }
@Table("PRODUCTS")
record ProductEntity(@Id String id, String name, int stock) {}
@Table("ORDERS")
record OrderEntity(@Id Integer id, String userId, String productId) {}
@Repository
static class R2dbcProductRepositoryAdapter implements ProductRepository {
private final R2dbcEntityTemplate template;
public R2dbcProductRepositoryAdapter(R2dbcEntityTemplate t) { this.template = t; }
public Mono<Product> findById(String id) { return template.selectOne(query(where("id").is(id)),ProductEntity.class).map(e -> new Product(e.id(), e.stock())); }
public Mono<Void> updateStock(String id, int stock) { return template.getDatabaseClient().sql("UPDATE PRODUCTS SET stock = :s WHERE id = :i").bind("s", stock).bind("i", id).fetch().rowsUpdated().then(); }
}
@Repository
static class R2dbcOrderRepositoryAdapter implements OrderRepository {
private final R2dbcEntityTemplate template;
public R2dbcOrderRepositoryAdapter(R2dbcEntityTemplate t) { this.template = t; }
public Mono<Order> save(Order o) { return template.insert(new OrderEntity(null, o.userId(), o.productId())).map(e -> o); }
}
@TransactionalWrapperUseCase
static class PurchaseProductUseCase implements PurchaseProductUseCasePort {
private final ProductRepository pRepo;
private final OrderRepository oRepo;
public PurchaseProductUseCase(ProductRepository p, OrderRepository o) { this.pRepo = p; this.oRepo = o; }
public Mono<Void> execute(String userId, String productId) {
return pRepo.findById(productId).flatMap(p -> pRepo.updateStock(p.id(), p.stock() - 1)).then(oRepo.save(new Order(userId, productId))).then();
}
}
}
Conclusión y Próximos Pasos
Hemos construido una solución completa, limpia y robusta para un problema complejo. Al mantener nuestro dominio puro y delegar las responsabilidades transversales a la capa de aplicación mediante AOP, logramos un código desacoplado, mantenible y altamente testeable. Las dependencias del proyecto reflejan esta arquitectura limpia, utilizando starters de Spring Boot para AOP y R2DBC, y librerías de prueba para H2 y ArchUnit.
// build.gradle
dependencies {
implementation 'org.reactivecommons.utils:object-mapper:0.1.0'
implementation project(':r2dbc-postgresql')
implementation project(':reactive-web')
implementation project(':model')
implementation project(':usecase')
implementation 'org.springframework.boot:spring-boot-starter'
implementation 'org.springframework.boot:spring-boot-starter-aop'
implementation 'org.springframework.boot:spring-boot-starter-data-r2dbc'
runtimeOnly('org.springframework.boot:spring-boot-devtools')
testImplementation 'com.tngtech.archunit:archunit:1.4.1'
testImplementation 'com.fasterxml.jackson.core:jackson-databind'
testImplementation 'com.h2database:h2'
testImplementation 'io.r2dbc:r2dbc-h2'
}
Este patrón no se limita a las transacciones. El mismo mecanismo de anotación y aspecto puede extenderse para manejar otras responsabilidades, como la autorización de seguridad, la auditoría o el registro de métricas, consolidándose como una base sólida para el desarrollo de futuras funcionalidades en cualquier aplicación reactiva que aspire a una arquitectura limpia y escalable.