Gestión de la Idempotencia en la Publicación de Productos para Venta Flash
La implementación de una funcionalidad de venta flash (seckill) en una plataforma de comercio electrónico requiere una cuidadosa consideración de la idempotencia, especialmente durante el proceso de publicación de productos. Este artículo detalla cómo abordar este desafío, asegurando que las operaciones repetidas no causen efectos secundarios no deseados.
Servicio de Venta Flash (SecKillServiceImpl)
El servicio central para la venta flash se encarga de recuperar y almacenar información relevante de productos y sesiones de venta flash en Redis. La lógica principal se enfoca en:
- Recuperación de Sesiones y Productos: Obtener las sesiones de venta flash activas para los próximos tres días a través de un srevicio de cupones (
CouponFeignService). - Almacenamiento de Información de Sesiones: Guardar los detalles de cada sesión de venta flash en Redis, utiilzando una clave que incluye el ID, nombre y rango de tiempo de la sesión.
- Almacenamiento de Información de Productos (SKUs): Para cada producto asociado a una sesión de venta flash:
- Generar un código aleatorio único (token) para identificar la instancia de la venta flash.
- Recuperar la información detallada del SKU del producto mediante un servicio de productos (
ProductFeignService). - Construir un objeto DTO (
SecKillSkuRedisTo) que contenga la información del SKU, los detalles de la promoción de venta flash y el código aleatorio. - Almacenar este objeto DTO serializado en JSON en una estructura hash de Redis, identificada por una clave combinada del ID de sesión y el ID del SKU.
- Control de Acceso Concurrente (Idempotnecia): Establecer un semáforo distribuido en Redis (usando Redisson) con el nombre basado en el código aleatorio y el prefijo de stock de SKU. El número de permisos inicial del semáforo se establece igual a la cantidad de unidades disponibles para la venta flash (
seckillCount). Esto previene la publicación duplicada del mismo lote de productos para venta flash.
Ejemplo de Código (Servicio de Venta Flash):
package com.alatus.mall.seckill.service.impl;
import com.alatus.common.utils.R; // Asumiendo una clase de respuesta genérica
import com.alatus.mall.seckill.constant.SecKillConstants;
import com.alatus.mall.seckill.feign.CouponFeignService;
import com.alatus.mall.seckill.feign.ProductFeignService;
import com.alatus.mall.seckill.service.SecKillService;
import com.alatus.mall.seckill.to.SecKillSkuRedisTo;
import com.alatus.mall.seckill.vo.SeckillSessionEntityVo;
import com.alatus.mall.seckill.vo.SkuInfoVo;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.TypeReference;
import org.redisson.api.RSemaphore;
import org.redisson.api.RedissonClient;
import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.BoundHashOperations;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Service;
import java.util.List;
import java.util.UUID;
import java.util.stream.Collectors;
@Service
public class SecKillServiceImpl implements SecKillService {
@Autowired
private CouponFeignService couponFeignService;
@Autowired
private ProductFeignService productFeignService;
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private RedissonClient redissonClient;
@Override
public void publishLatest3DaysSecKillProducts() {
// Recuperar sesiones de venta flash activas
R activeSessionsResponse = couponFeignService.fetchActiveSessionsIn3Days();
if (activeSessionsResponse.getCode() == 0) { // Asumiendo 0 para éxito
List<seckillsessionentityvo> activeSessions = activeSessionsResponse.getData(new TypeReference<list>>() {});
cacheSessionDetails(activeSessions);
cacheProductDetailsForSessions(activeSessions);
}
}
private void cacheSessionDetails(List<seckillsessionentityvo> sessions) {
sessions.forEach(session -> {
long startTimeMillis = session.getStartTime().getTime();
long endTimeMillis = session.getEndTime().getTime();
// Clave única para la sesión: ID + prefijo + Nombre + Rango de tiempo
String sessionCacheKey = String.format("%d_%s_%d-%d",
session.getId(),
SecKillConstants.SESSION_CACHE_PREFIX,
startTimeMillis,
endTimeMillis);
// Almacenar la lista de SKUs en esta sesión si la clave no existe
if (!redisTemplate.hasKey(sessionCacheKey)) {
List<string> skuSessionIds = session.getRelationSkus().stream()
.map(item -> String.format("%d_%d", item.getPromotionSessionId(), item.getSkuId()))
.collect(Collectors.toList());
redisTemplate.opsForList().leftPushAll(sessionCacheKey, skuSessionIds);
}
});
}
private void cacheProductDetailsForSessions(List<seckillsessionentityvo> sessions) {
sessions.forEach(session -> {
// Operación de hash en Redis para almacenar detalles de SKUs de venta flash
BoundHashOperations<string string=""> skuHashOps = redisTemplate.boundHashOps(SecKillConstants.SKU_CACHE_PREFIX);
session.getRelationSkus().forEach(promotionSku -> {
// Generar un identificador único para esta instancia de venta flash
String uniqueFlashToken = UUID.randomUUID().toString().replace("-", "");
// Clave para identificar el SKU específico dentro de la sesión
String skuSessionKey = String.format("%d_%d", promotionSku.getPromotionSessionId(), promotionSku.getSkuId());
// Proceder solo si este SKU no ha sido cacheado previamente para esta sesión
if (!skuHashOps.hasKey(skuSessionKey)) {
SecKillSkuRedisTo seckillSkuData = new SecKillSkuRedisTo();
// Obtener información básica del SKU
R skuInfoResponse = productFeignService.fetchSkuInfo(promotionSku.getSkuId());
if (skuInfoResponse.getCode() == 0) {
SkuInfoVo skuDetails = skuInfoResponse.get("skuInfo", new TypeReference<skuinfovo>() {});
seckillSkuData.setSkuInfo(skuDetails);
}
// Copiar detalles de la promoción de venta flash
BeanUtils.copyProperties(promotionSku, seckillSkuData);
// Establecer tiempos de inicio y fin de la venta flash
seckillSkuData.setStartTime(session.getStartTime().getTime());
seckillSkuData.setEndTime(session.getEndTime().getTime());
seckillSkuData.setRandomCode(uniqueFlashToken);
// Serializar y almacenar en Redis
String skuJson = JSON.toJSONString(seckillSkuData);
skuHashOps.put(skuSessionKey, skuJson);
// Configurar el semáforo distribuido para el control de stock y concurrencia
String semaphoreKey = SecKillConstants.SKU_STOCK_SEMAPHORE + uniqueFlashToken;
RSemaphore stockSemaphore = redissonClient.getSemaphore(semaphoreKey);
// Inicializar el semáforo con la cantidad de stock disponible para la venta flash
// 'trySetPermits' asegura que si ya existe, no se modifica, lo que contribuye a la idempotencia.
stockSemaphore.trySetPermits(promotionSku.getSeckillCount());
}
});
});
}
// Método adicional para obtener los SKUs de venta flash actuales (ejemplo)
@Override
public List<seckillskuredisto> getCurrentSecKillSkus() {
// Implementación para recuperar SKUs de venta flash activos de Redis
// ...
return null; // Placeholder
}
}
</seckillskuredisto></skuinfovo></string></seckillsessionentityvo></string></seckillsessionentityvo></list></seckillsessionentityvo>
Tarea Programada para Publicación (SecKillSkuScheduled)
Una tarea programada se encarga de ejecutar la lógica de publicación de productos de venta flash de manera regular. Para garantizar la idempotencia a nivel de la tarea programada y evitar ejecuciones simultáneas del mismo trabajo, se utiliza un bloqueo distribuido implementado con Redisson:
- Bloqueo Distribuido: Se adquiere un bloqueo (
RLock) antes de ejecutar el servicio de publicación. Este bloqueo tiene una duración de 30 segundos. - Ejecución Idempotente: El servicio
SecKillService.publishLatest3DaysSecKillProducts()se llama dentro del bloquetrydel bloqueo. Si la tarea se ejecuta varias veces en rápida sucesión (por ejemplo, debido a reinicios del servidor o fallos temporales), solo una instancia adquirirá el bloqueo y procederá. Las otras esperarán o fallarán si el bloqueo expira. - Liberación del Bloqueo: El bloqueo se libera en el bloque
finallypara asegurar que esté disponible para futuras ejecuciones.
Ejemplo de Código (Tarea Programada):
package com.alatus.mall.seckill.scheduled;
import com.alatus.mall.seckill.service.SecKillService;
import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component; // Usar Component en lugar de Service si solo es un bean de programación
import java.util.concurrent.TimeUnit;
/**
* Programación para la publicación diaria de productos de venta flash.
* Publica productos de venta flash para los próximos tres días cada día a las 4 AM.
* Maneja la publicación duplicada de forma idempotente.
*/
@Component
@Slf4j
public class SecKillSkuScheduled {
@Autowired
private SecKillService secKillService;
@Autowired
private RedissonClient redissonClient;
// Clave para el bloqueo distribuido de la tarea de publicación
private static final String PUBLISH_UPLOAD_LOCK_KEY = "seckill:publish:upload:lock";
// Programación: Cada día a las 4:00 AM (formato cron)
// @Scheduled(cron = "0 0 4 * * ?")
// Programación de ejemplo para pruebas: cada 15 segundos
@Scheduled(cron = "*/15 * * * * ?")
public void schedulePublishLatest3DaysSecKillProducts() {
log.info("Iniciando tarea programada para publicar productos de venta flash...");
RLock publishLock = redissonClient.getLock(PUBLISH_UPLOAD_LOCK_KEY);
// Intentar adquirir el bloqueo con un tiempo de espera y una duración de lease
boolean acquired = false;
try {
// Espera hasta 30 segundos para adquirir el bloqueo, y el bloqueo durará 60 segundos si se adquiere
acquired = publishLock.tryLock(30, 60, TimeUnit.SECONDS);
if (acquired) {
log.info("Bloqueo adquirido. Ejecutando la publicación de productos de venta flash.");
secKillService.publishLatest3DaysSecKillProducts();
log.info("Publicación de productos de venta flash completada.");
} else {
log.warn("No se pudo adquirir el bloqueo para la publicación de productos de venta flash. Otra instancia podría estar ejecutándose.");
}
} catch (InterruptedException e) {
log.error("La tarea de publicación de venta flash fue interrumpida.", e);
Thread.currentThread().interrupt(); // Restaurar el estado de interrupción
} catch (Exception e) {
log.error("Ocurrió un error durante la ejecución de la publicación de productos de venta flash.", e);
} finally {
if (acquired && publishLock.isHeldByCurrentThread()) {
publishLock.unlock();
log.info("Bloqueo liberado.");
}
}
}
}
Controlador de API (SeckillController)
Se proporciona un endpoint de API para recuperar la lista de productos de venta flash actualmente disponibles para que los clientes interactúen con ellos.
Ejemplo de Código (Controlador):
package com.alatus.mall.seckill.app;
import com.alatus.common.utils.R; // Asumiendo una clase de respuesta genérica
import com.alatus.mall.seckill.service.SecKillService;
import com.alatus.mall.seckill.to.SecKillSkuRedisTo;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import java.util.List;
@RestController
@RequestMapping("/api/seckill") // Definir un prefijo base para la API
public class SeckillController {
@Autowired
private SecKillService secKillService;
/**
* Obtiene la lista de productos de venta flash actualmente activos.
* @return Una respuesta R que contiene la lista de productos de venta flash.
*/
@GetMapping("/current-products")
public R getCurrentActiveSeckillProducts() {
List<seckillskuredisto> activeSkus = secKillService.getCurrentSecKillSkus();
if (activeSkus != null) {
return R.ok("Se recuperaron los productos de venta flash activos.").put("data", activeSkus);
} else {
// Considerar un caso donde no hay productos activos o hay un error al recuperarlos.
return R.ok("No hay productos de venta flash activos en este momento.").put("data", List.of());
}
}
}
</seckillskuredisto>
Al combinar el uso de semáforos para el control de stock a nivel de SKU y bloqueos distribuidos para la tarea de publicación, el sistema logra una solución robusta y idempotente para la gestión de la publicación de productos de venta flash.