برمجة Spring WebFl التفاعلية
WebFlux هو محرك Spring غير المتزامن—بفضل Mono الفردي وFlux من الدرجة الأولى، يتعامل الإدخال/الإخراج غير المحظور مع عشرات الآلاف من الاتصالات المتزامنة.
1. ما ستتعلمه
- نموذج البرمجة التفاعلية: المفاهيم الأساسية وعوامل Mono / Flux
@RestController+Mono<T>/Flux<T>كتابة واجهات برمجة التطبيقات التفاعلية- WebFlux مقابل WebMVC: مقارنة حالات الاستخدام ودليل الاختيار
- المستودعات التفاعلية (Spring Data R2DBC)
- Alice تنفذ واجهة برمجة تطبيقات تفاعلية قادرة على معالجة آلاف استعلامات الطلبات في الثانية
2. قصة حقيقية من مهندس التزامن العالي
(1) نقطة الألم: استنفاد تجمع الخيوط
واجهت OrderFlow الخاصة بـ Alice اختناقًا أثناء حدث مبيعات كبير—مع 500 مستخدم متزامن يستعلون عن الطلبات، تم استنفاد تجمع خيوط Tomcat (المعين إلى 200 خيط افتراضيًا)، وتم وضع الطلبات في قائمة الانتظار، وقفز زمن الاستجابة P99 إلى 3 ثوانٍ. معظم الوقت، كانت الخيوط تنتظر الإدخال/الإخراج (قاعدة بيانات، شبكة)، وكان استخدام المعالج 15% فقط. اقترح Bob زيادة عدد الخيوط، لكن إضافة 1,000 خيط ستستهلك 1 غيغابايت إضافي من الذاكرة.
(2) حل WebFlux
يتعامل WebFlux مع التزامن العالي بعدد قليل من الخيوط—لا تُحظر الخيوط أثناء انتظار الإدخال/الإخراج ويمكنها الاستمرار في خدمة الطلبات الأخرى:
@GetMapping("/{id}")
public Mono<OrderResponse> getOrder(@PathVariable Long id) {
return orderRepository.findById(id)
.map(OrderResponse::from)
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Order", id)));
}
(3) الإيرادات
بعد أن أعادت Alice بناء واجهة برمجة استعلام الطلبات باستخدام WebFlux، أصبحت أربعة خيوط قادرة على التعامل مع 2,000 طلب متزامن، وانخفض استخدام الذاكرة بنسبة 70%، وانخفض زمن الاستجابة P99 من 3 ثوانٍ إلى 100 مللي ثانية، وارتفع استخدام المعالج إلى 60%.
3. المفاهيم الأساسية للبرمجة التفاعلية
(1) Mono وFlux
| النوع | المعنى | التشبيه | مثال |
|---|---|---|---|
Mono<T> |
0 أو 1 عنصر | نسخة غير متزامنة من Optional | استعلام طلب واحد |
Flux<T> |
من 0 إلى N عنصر | نسخة غير متزامنة من Stream | استعلام قائمة الطلبات |
Mono<Void> |
لا قيمة إرجاع | إشارة اكتمال | عملية حذف |
graph LR
A["Publisher"] --> B["Mono<br/>0 أو 1 عنصر"]
A --> C["Flux<br/>من 0 إلى N عنصر"]
B --> D["onNext → onComplete"]
C --> E["onNext x N → onComplete"]
B --> F["onError"]
C --> F
(1) ▶ مثال: العمليات الأساسية في Mono وFlux
// Mono: قيمة واحدة
Mono<String> mono = Mono.just("Hello OrderFlow");
Mono<String> empty = Mono.empty();
Mono<String> fromCallable = Mono.fromCallable(() -> fetchOrder(1L));
// Flux: قيم متعددة
Flux<Integer> flux = Flux.just(1, 2, 3, 4, 5);
Flux<Integer> range = Flux.range(1, 100);
Flux<Long> interval = Flux.interval(Duration.ofSeconds(1));
الناتج:
// التنفيذ ناجح
(2) العوامل الشائعة
| العامل | الوظيفة | مثال |
|---|---|---|
map |
تحويل العناصر | .map(OrderResponse::from) |
flatMap |
تحويل غير متزامن | .flatMap(this::enrichOrder) |
filter |
تصفية | .filter(order -> "PENDING".equals(order.getStatus())) |
switchIfEmpty |
بديل عند القيمة الفارغة | .switchIfEmpty(Mono.error(...)) |
onErrorResume |
تخفيض في حالة الخطأ | .onErrorResume(e -> fallbackOrder()) |
timeout |
التحكم في المهلة | .timeout(Duration.ofSeconds(5)) |
retry |
إعادة المحاولة | .retry(3) |
subscribeOn |
تحديد خيط الاشتراك | .subscribeOn(Schedulers.boundedElastic()) |
(2) ▶ مثال: تسلسل العوامل
public Mono<OrderResponse> getOrderWithDetails(Long id) {
return orderRepository.findById(id)
.flatMap(order ->
productRepository.findById(order.getProductId())
.map(product -> OrderResponse.from(order, product)))
.switchIfEmpty(Mono.error(
new ResourceNotFoundException("Order", id)))
.timeout(Duration.ofSeconds(3))
.onErrorResume(TimeoutException.class,
e -> Mono.error(new RuntimeException("انتهت مهلة استعلام الطلب")))
.retry(2);
}
الناتج:
// التنفيذ ناجح
4. واجهة برمجة التطبيقات REST في WebFlux
(1) ▶ مثال: المتحكم التفاعلي
@RestController
@RequestMapping("/api/v1/reactive/orders")
public class ReactiveOrderController {
private final ReactiveOrderService orderService;
public ReactiveOrderController(ReactiveOrderService orderService) {
this.orderService = orderService;
}
@GetMapping("/{id}")
public Mono<OrderResponse> getOrder(@PathVariable Long id) {
return orderService.findById(id);
}
@GetMapping
public Flux<OrderResponse> listOrders(
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "20") int size) {
return orderService.findAll(page, size);
}
@PostMapping
public Mono<ResponseEntity<OrderResponse>> createOrder(
@Valid @RequestBody CreateOrderRequest request) {
return orderService.create(request)
.map(order -> ResponseEntity
.status(HttpStatus.CREATED).body(order));
}
@DeleteMapping("/{id}")
public Mono<ResponseEntity<Void>> cancelOrder(@PathVariable Long id) {
return orderService.cancel(id)
.then(Mono.just(ResponseEntity.noContent().<Void>build()));
}
}
الناتج:
// التنفيذ ناجح
| البُعد | WebMVC | WebFlux |
|---|---|---|
| نوع الإرجاع | كائن متزامن | Mono<T> / Flux<T> |
| نموذج الخيوط | خيط واحد لكل طلب | حلقة الأحداث (خيوط قليلة) |
| نموذج الإدخال/الإخراج | محظور | غير محظور |
| الحاوية | Tomcat / Jetty | Netty |
| حد التزامن | حجم تجمع الخيوط | غير محدود تقريبًا |
5. الاختيار بين WebFlux وWebMVC
(1) قرار اختيار المنتج
| السيناريو | التوصية | السبب |
|---|---|---|
| واجهة برمجة تطبيقات CRUD، تزامن < 500 | WebMVC | بسيط وبديهي، نظام بيئي ناضج |
| كثيف الإدخال/الإخراج، تزامن عالي | WebFlux | غير محظور، يدعم التزامن العالي بعدد قليل من الخيوط |
| تدفق في الوقت الفعلي (SSE/WebSocket) | WebFlux | دعم أصلي |
| عدد كبير من العمليات المحظورة (JDBC) | WebMVC | العمليات المحظورة أسوأ فعليًا في WebFlux |
| سيناريوهات مختلطة | WebMVC + جزء غير متزامن | تحسين تدريجي |
(1) ▶ مثال: أحداث إرسال الخادم (SSE)
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<OrderEvent> streamOrderEvents() {
return orderEventPublisher.eventStream()
.log("order-events");
}
الناتج:
// التنفيذ ناجح
6. Spring Data R2DBC
(1) الوصول التفاعلي إلى قاعدة البيانات
(1) ▶ مثال: كيان ومستودع R2DBC
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-r2dbc</artifactId>
</dependency>
<dependency>
<groupId>io.r2dbc</groupId>
<artifactId>r2dbc-mysql</artifactId>
</dependency>
الناتج:
// التنفيذ ناجح
@Table("products")
public class Product {
@Id
private Long id;
private String name;
private BigDecimal price;
private Integer stock;
public Product() {}
public Product(String name, BigDecimal price, Integer stock) {
this.name = name;
this.price = price;
this.stock = stock;
}
// getters
}
public interface ReactiveProductRepository
extends ReactiveCrudRepository<Product, Long> {
Flux<Product> findByNameContaining(String keyword);
Flux<Product> findByStockLessThan(Integer threshold);
}
| البُعد | JPA + JDBC | R2DBC |
|---|---|---|
| نموذج المشغل | محظور | غير محظور |
| نمط واجهة برمجة التطبيقات | متزامن | تفاعلي (Mono/Flux) |
| تعيين العلاقات | مدعوم (@OneToMany) | غير مدعوم |
| لغة الاستعلام | JPQL | SQL أصلي / اشتقاق اسم الطريقة |
| المعاملات | @Transactional |
@Transactional (تفاعلي) |
(2) ▶ مثال: خدمة تفاعلية
@Service
public class ReactiveOrderService {
private final ReactiveOrderRepository orderRepository;
private final ReactiveProductRepository productRepository;
public ReactiveOrderService(ReactiveOrderRepository orderRepository,
ReactiveProductRepository productRepository) {
this.orderRepository = orderRepository;
this.productRepository = productRepository;
}
public Mono<OrderResponse> create(CreateOrderRequest request) {
return productRepository.findById(request.productId())
.switchIfEmpty(Mono.error(
new ResourceNotFoundException("Product", request.productId())))
.flatMap(product -> {
if (product.getStock() < request.quantity()) {
return Mono.error(new InsufficientStockException(
product.getId(), product.getStock(), request.quantity()));
}
product.deductStock(request.quantity());
return productRepository.save(product)
.then(orderRepository.save(new Order(product.getId(), request.quantity())));
})
.map(OrderResponse::from);
}
public Flux<OrderResponse> findAll(int page, int size) {
return orderRepository.findAll()
.skip(page * size)
.take(size)
.map(OrderResponse::from);
}
public Mono<Void> cancel(Long orderId) {
return orderRepository.findById(orderId)
.flatMap(order -> {
order.setStatus("CANCELLED");
return orderRepository.save(order);
})
.then();
}
}
الناتج:
// التنفيذ ناجح
7. مثال شامل: واجهة برمجة استعلام الطلبات التفاعلية OrderFlow
// ReactiveProductRepository.java
public interface ReactiveProductRepository
extends ReactiveCrudRepository<Product, Long> {
Flux<Product> findByNameContaining(String keyword);
}
// ReactiveOrderRepository.java
public interface ReactiveOrderRepository
extends ReactiveCrudRepository<Order, Long> {
Flux<Order> findByStatus(String status);
}
// ReactiveOrderController.java
@RestController
@RequestMapping("/api/v1/reactive/orders")
public class ReactiveOrderController {
private final ReactiveOrderService orderService;
public ReactiveOrderController(ReactiveOrderService orderService) {
this.orderService = orderService;
}
@GetMapping("/{id}")
public Mono<OrderResponse> get(@PathVariable Long id) {
return orderService.findById(id);
}
@GetMapping
public Flux<OrderResponse> list(
@RequestParam(defaultValue = "0") int page,
@RequestParam(defaultValue = "20") int size) {
return orderService.findAll(page, size);
}
@GetMapping("/status/{status}")
public Flux<OrderResponse> byStatus(@PathVariable String status) {
return orderService.findByStatus(status);
}
@PostMapping
public Mono<ResponseEntity<OrderResponse>> create(
@RequestBody CreateOrderRequest req) {
return orderService.create(req)
.map(o -> ResponseEntity.status(HttpStatus.CREATED).body(o));
}
@DeleteMapping("/{id}")
public Mono<ResponseEntity<Void>> cancel(@PathVariable Long id) {
return orderService.cancel(id)
.thenReturn(ResponseEntity.noContent().<Void>build());
}
}
// ReactiveOrderService.java
@Service
public class ReactiveOrderService {
private final ReactiveOrderRepository orderRepo;
private final ReactiveProductRepository productRepo;
public ReactiveOrderService(ReactiveOrderRepository orderRepo,
ReactiveProductRepository productRepo) {
this.orderRepo = orderRepo;
this.productRepo = productRepo;
}
public Mono<OrderResponse> findById(Long id) {
return orderRepo.findById(id)
.map(OrderResponse::from)
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Order", id)));
}
public Flux<OrderResponse> findAll(int page, int size) {
return orderRepo.findAll().skip((long) page * size).take(size)
.map(OrderResponse::from);
}
public Flux<OrderResponse> findByStatus(String status) {
return orderRepo.findByStatus(status).map(OrderResponse::from);
}
public Mono<OrderResponse> create(CreateOrderRequest req) {
return productRepo.findById(req.productId())
.switchIfEmpty(Mono.error(new ResourceNotFoundException("Product", req.productId())))
.flatMap(p -> {
if (p.getStock() < req.quantity()) {
return Mono.error(new InsufficientStockException(p.getId(), p.getStock(), req.quantity()));
}
p.deductStock(req.quantity());
return productRepo.save(p)
.then(orderRepo.save(new Order(p.getId(), req.quantity())));
})
.map(OrderResponse::from);
}
public Mono<Void> cancel(Long id) {
return orderRepo.findById(id)
.flatMap(o -> { o.setStatus("CANCELLED"); return orderRepo.save(o); })
.then();
}
}
❓ أسئلة شائعة
flatMap.Mono.fromCallable() + subscribeOn(Schedulers.boundedElastic()) لجدولة العمليات المحظورة إلى تجمع الخيوط المرن. ومع ذلك، الاستخدام المتكرر يشير إلى أن هذا النهج غير مناسب لـ WebFlux.onErrorResume / onErrorReturn؛ 2) @ExceptionHandler العام (مدعوم أيضًا من WebFlux)؛ 3) onErrorMap لتحويل أنواع الاستثناءات.WebTestClient بدلاً من MockMvc. يقوم @WebFluxTest بتحميل طبقة WebFlux. يُستخدم StepVerifier لاختبار تسلسلات العناصر في Mono/Flux.📖 ملخص
- Mono (0/1 عنصر) وFlux (من 0 إلى N عنصر) هما النوعان الأساسيان في البرمجة التفاعلية
- تسلسل استدعاء العوامل: map، flatMap، filter، switchIfEmpty، timeout، retry
- يستخدم WebFlux Netty وحلقة الأحداث لدعم التزامن العالي بعدد صغير من الخيوط
- WebFlux مناسب لسيناريوهات الإدخال/الإخراج المكثف والتزامن العالي، لكنه غير مناسب لعدد كبير من العمليات المحظورة
- R2DBC هو مشغل قاعدة بيانات تفاعلي؛ لا يدعم تعيين العلاقات، لذا يجب دمج البيانات يدويًا باستخدام
flatMap - يجب عدم استخدام WebMVC وWebFlux معًا؛ اختر أحدهما بناءً على السيناريو المحدد
📝 تمارين
-
مسألة أساسية (الصعوبة ⭐): نفذ واجهة برمجة تطبيقات تفاعلية لاستعلام المنتجات باستخدام WebFlux (
GET /api/v1/reactive/products/{id})، مع الاتصال بقاعدة بيانات H2 في الذاكرة عبر R2DBC. -
تمرين متقدم (الصعوبة: ⭐⭐): نفذ واجهة برمجة تطبيقات تفاعلية لوضع الطلبات تغطي العملية بالكامل—بما في ذلك فحص المخزون وخصم المخزون وإنشاء الطلب—باستخدام
flatMapلدمج عمليات تفاعلية متعددة. -
تحدٍ (الصعوبة: ⭐⭐⭐): استخدم
WebTestClientوStepVerifierلكتابة اختبار لواجهة برمجة تطبيقات تفاعلية، ونفذ واجهة دفع أحداث الطلبات في الوقت الفعلي SSE، وقارن اختلافات الأداء بين WebMVC وWebFlux تحت 1,000 اتصال متزامن.



