ترجمهٔ مدل (Model Translation)
زمینهٔ محدود مرز یک مدل، یعنی زبان فراگیر، است. همانطور که در فصل ۳ آموختید، الگوهای مختلفی برای طراحی ارتباط میان زمینههای محدود وجود دارد. اگر تیمهای دو زمینهٔ محدود ارتباط مؤثر و تمایل به همکاری دارند، میتوان Partnership داشت: پروتکلها بهصورت موردی هماهنگ میشوند و مشکلات یکپارچهسازی از طریق ارتباط تیمها حل میشوند. روش همکاریمحور دیگر Shared Kernel است که در آن تیمها بخش محدودی از مدل، مثلاً قراردادهای Integration، را استخراج و مشترکاً تکامل میدهند.
در رابطهٔ Customer–Supplier، قدرت به سمت Upstream یا Downstream متمایل است. اگر Downstream نتواند با مدل Upstream همرنگ شود، راهحل فنی پیچیدهتری برای ترجمهٔ مدلها لازم است.
این ترجمه میتواند توسط یکی یا هر دو طرف انجام شود: Downstream با Anticorruption Layer (ACL) مدل Upstream را به نیاز خودش تطبیق میدهد و Upstream با Open-Host Service (OHS) و Published Language مخصوص Integration، مصرفکنندگان را از تغییرات مدل داخلی محافظت میکند. چون منطق ترجمه در ACL و OHS مشابه است، این فصل گزینههای پیادهسازی را بدون جداسازی مصنوعی دو الگو بررسی میکند و فقط در موارد لازم تفاوت را یادآور میشود.
ترجمهٔ مدل میتواند Stateless یا Stateful باشد. Stateless هنگام رسیدن درخواست ورودی OHS یا صدور درخواست خروجی ACL و در لحظه انجام میشود. Stateful به منطق پیچیدهتر و Database نیاز دارد.
ترجمهٔ Stateless
زمینهای که مالک ترجمه است ــ OHS در Upstream یا ACL در Downstream ــ الگوی Proxy را برای Intercept کردن Requestهای ورودی/خروجی و Map کردن Source Model به Target Model پیاده میکند؛ شکل ۹-۱.
شکل ۹-۱ — ترجمهٔ مدل توسط Proxy.
پیادهسازی Proxy به Synchronous یا Asynchronous بودن ارتباط بستگی دارد.
همگام (Synchronous)
روش معمول، قراردادن منطق Transformation در Codebase همان Bounded Context است؛ شکل ۹-۲. در OHS ترجمه به Public Language هنگام پردازش Request ورودی و در ACL هنگام فراخوانی Upstream انجام میشود.
شکل ۹-۲ — ارتباط Synchronous.
در برخی موارد، سپردن ترجمه به جزء بیرونی مانند API Gateway اقتصادیتر و راحتتر است. Gateway میتواند محصول Open Source مانند Kong یا KrakenD یا Managed Service ابری مانند AWS API Gateway، Google Apigee یا Azure API Management باشد.
برای OHS، API Gateway مدل داخلی را به Published Language بهینهشده برای Integration تبدیل میکند و مدیریت و ارائهٔ چند Version از API را نیز آسان میسازد؛ شکل ۹-۳.
شکل ۹-۳ — ارائهٔ Versionهای مختلف Published Language.
ACL پیادهشده با API Gateway میتواند توسط چند Downstream Context مصرف شود. در این حالت ACL عملاً یک Bounded Context مخصوص Integration است؛ شکل ۹-۴. چنین Contextهایی که مسئولیت اصلیشان تبدیل مدل برای مصرف آسانتر سایر اجزاست، Interchange Context نامیده میشوند.
شکل ۹-۴ — ACL مشترک.
ناهمگام (Asynchronous)
برای ترجمهٔ مدل در ارتباط ناهمگام میتوان Message Proxy ساخت: جزء میانی به پیامهای Source Context Subscribe میکند، Transformation لازم را اعمال و پیام حاصل را به Target Subscriber Forward میکند؛ شکل ۹-۵.
شکل ۹-۵ — ترجمهٔ مدل در ارتباط Asynchronous.
این جزء علاوه بر ترجمه میتواند پیامهای نامرتبط را Filter کند و Noise واردشده به Target را کاهش دهد.
ترجمهٔ Asynchronous برای OHS حیاتی است. اشتباه رایج این است که برای Objectهای مدل Published Language طراحی شود اما Domain Eventها بدون ترجمه منتشر شوند؛ در این صورت Implementation Model لو میرود. Proxy میتواند Domain Eventها را Intercept و به Published Language تبدیل کند؛ شکل ۹-۶.
این کار همچنین تفاوت میان Private Eventهای مخصوص نیازهای داخلی Context و Public Eventهای طراحیشده برای Integration را روشن میکند. فصل ۱۵ این موضوع را در رابطهٔ DDD و Event-Driven Architecture بسط میدهد.
شکل ۹-۶ — Domain Eventها در Published Language.
ترجمهٔ Stateful
برای Transformationهای عمدهتر ــ مثلاً وقتی باید دادهٔ Source تجمیع شود یا دادهٔ چند Source در یک Model واحد ترکیب گردد ــ ترجمهٔ Stateful لازم است.
تجمیع دادهٔ ورودی
فرض کنید Bounded Context برای Performance میخواهد Requestهای ورودی را جمع و Batch Processing کند. این Aggregation هم برای Requestهای همگام و هم ناهمگام ممکن است لازم باشد؛ شکل ۹-۷.
شکل ۹-۷ — Batch کردن Requestها.
کاربرد رایج دیگر، ترکیب چند پیام ریزدانه در یک پیام واحد با دادهٔ یکپارچه است؛ شکل ۹-۸.
شکل ۹-۸ — یکپارچهکردن Eventهای ورودی.
Aggregation منبع را نمیتوان با API Gateway ساده انجام داد و به پردازش Stateful نیاز دارد. منطق ترجمه برای Track کردن دادهٔ ورودی و پردازش بعدی باید Storage Persistشوندهٔ خودش را داشته باشد؛ شکل ۹-۹.
شکل ۹-۹ — Stateful Model Transformation.
در بعضی موارد میتوان بهجای ساخت Solution اختصاصی از محصول آماده استفاده کرد؛ مثلاً Stream Processing با Kafka یا AWS Kinesis، یا Batch Processing با Apache NiFi، AWS Glue یا Spark.
یکپارچهکردن چند Source
یک Bounded Context ممکن است به Aggregate کردن داده از چند Source، شامل Contextهای دیگر، نیاز داشته باشد. نمونهٔ معمول Backend-for-Frontend است که UI باید دادهٔ چند Service را ترکیب کند.
نمونهٔ دیگر، Contextی است که دادهٔ چند Context دیگر را میگیرد و Business Logic پیچیدهای روی آنها اجرا میکند. در این حالت میتوان پیچیدگی Integration را از منطق جدا کرد و ACLای در جلوی Context قرار داد که دادهٔ همهٔ Sourceها را تجمیع کند؛ شکل ۹-۱۰.
شکل ۹-۱۰ — سادهسازی Integration Model با ACL.
یکپارچهسازی Aggregateها
در فصل ۶ دیدیم Aggregateها از طریق Domain Event با بقیهٔ سیستم ارتباط برقرار میکنند و اجزای بیرونی میتوانند Subscribe کنند. اما Domain Event چگونه به Message Bus منتشر میشود؟ پیش از پاسخ، چند پیادهسازی اشتباه رایج را بررسی کنیم.
کد زیر Event را داخل خود Aggregate منتشر میکند:
01 public class Campaign
02 {
03 ...
04 List<DomainEvent> _events;
05 IMessageBus _messageBus;
06 ...
07
08 public void Deactivate(string reason)
09 {
10 for (l in _locations.Values())
11 {
12 l.Deactivate();
13 }
14
15 IsActive = false;
16
17 var newEvent = new CampaignDeactivated(_id, reason);
18 _events.Append(newEvent);
19 _messageBus.Publish(newEvent);
20 }
21 }
این پیادهسازی ساده اما غلط است. Event پیش از Commit شدن State جدید Aggregate منتشر میشود؛ Subscriber ممکن است اعلان Deactivate را بگیرد در حالی که Database هنوز State قبلی را دارد. دوم، اگر Commit بعداً بهدلیل Race Condition، منطق بعدی یا خطای فنی شکست بخورد، Transaction Rollback میشود اما Event قبلاً منتشر شده و قابل پسگرفتن نیست.
روش دوم انتشار را به Application Layer منتقل میکند و ابتدا Aggregate را Commit میکند:
01 public class ManagementAPI
02 {
03 ...
04 private readonly IMessageBus _messageBus;
05 private readonly ICampaignRepository _repository;
06 ...
07 public ExecutionResult DeactivateCampaign(CampaignId id, string reason)
08 {
09 try
10 {
11 var campaign = repository.Load(id);
12 campaign.Deactivate(reason);
13 _repository.CommitChanges(campaign);
14
15 var events = campaign.GetUnpublishedEvents();
16 for (IDomainEvent e in events)
17 {
18 _messageBus.publish(e);
19 }
20 campaign.ClearUnpublishedEvents();
21 }
22 catch(Exception ex)
23 {
24 ...
25 }
26 }
27 }
این نیز کاملاً قابل اعتماد نیست. اگر Process پس از Commit Database و پیش از Publish Event Crash کند یا Message Bus در دسترس نباشد، State Commit شده ولی Event هرگز منتشر نمیشود. برای این Edge Case از Outbox استفاده میکنیم.
Outbox
الگوی Outbox، انتشار قابل اتکای Domain Eventها را با الگوریتم زیر تضمین میکند؛ شکل ۹-۱۱:
- State جدید Aggregate و Domain Eventهای جدید در همان تراکنش اتمیک Commit میشوند.
- Message Relay رویدادهای تازهCommitشده را از Database میخواند.
- Relay آنها را روی Message Bus منتشر میکند.
- پس از Publish موفق، رویداد در Database Published علامتگذاری یا حذف میشود.
شکل ۹-۱۱ — الگوی Outbox.
در Relational Database میتوان State و جدول اختصاصی Outbox را اتمیک Commit کرد؛ شکل ۹-۱۲.
شکل ۹-۱۲ — جدول Outbox.
در NoSQL که Multi-document Transaction ندارد، Eventهای خروجی باید داخل رکورد خود Aggregate Embed شوند:
{
"campaign-id": "364b33c3-2171-446d-b652-8e5a7b2be1af",
"state": {
"name": "Autumn 2017",
"publishing-state": "DEACTIVATED",
"ad-locations": [ ... ]
...
},
"outbox": [
{
"campaign-id": "364b33c3-2171-446d-b652-8e5a7b2be1af",
"type": "campaign-deactivated",
"reason": "Goals met",
"published": false
}
]
}
Fetch کردن Eventهای منتشرنشده
Relay میتواند به دو روش کار کند:
- Pull / Polling Publisher: Database را پیوسته برای Eventهای Publishedنشده Query کند. Index مناسب برای کاهش Load ضروری است.
- Push / Transaction Log Tailing: از قابلیت Database برای Notification تغییرات استفاده کند؛ مثلاً دنبالکردن Transaction Log در بعضی Relational Databaseها یا Stream تغییرات در DynamoDB Streams.
Outbox تحویل پیام را حداقل یکبار (At Least Once) تضمین میکند. اگر Relay پس از Publish و پیش از Mark کردن پیام Crash کند، پیام در Iteration بعدی دوباره Publish میشود؛ پس مصرفکنندگان باید Duplicate را تحمل کنند.
Saga
یک اصل Aggregate این است که هر Transaction فقط یک Aggregate Instance را تغییر دهد. اما بعضی Business Processها چند Aggregate را دربر میگیرند. مثلاً هنگام Activate شدن یک کمپین تبلیغاتی، مواد تبلیغ باید برای Publisher ارسال شوند؛ در صورت تأیید، State کمپین Published شود و در صورت ردشدن، Rejected.
این جریان دو Business Entity ــ Campaign و Publisher ــ را درگیر میکند که مسئولیتهای متفاوت و شاید Bounded Contextهای جدا دارند. قراردادن آنها در یک Aggregate بیشمهندسی است. این جریان را میتوان با Saga پیاده کرد.
Saga یک Business Process بلندمدت است؛ نه لزوماً از نظر زمان ــ ممکن است ثانیهها یا سالها طول بکشد ــ بلکه از نظر تعداد Transactionها. Saga به Eventهای اجزای مربوط گوش میدهد و Commandهای بعدی را صادر میکند. اگر مرحلهای شکست بخورد، Saga Action جبرانی مناسب را صادر میکند تا State سیستم سازگار بماند.
در مثال کمپین، Saga باید به CampaignActivated و PublishingConfirmed/PublishingRejected گوش دهد و SubmitAdvertisement و TrackPublishingConfirmation/TrackPublishingRejection را اجرا کند؛ شکل ۹-۱۳. TrackPublishingRejection در اینجا Compensation است تا کمپین بهاشتباه Active باقی نماند.
شکل ۹-۱۳ — Saga.
public class CampaignPublishingSaga
{
private readonly ICampaignRepository _repository;
private readonly IPublishingServiceClient _publishingService;
...
public void Process(CampaignActivated @event)
{
var campaign = _repository.Load(@event.CampaignId);
var advertisingMaterials = campaign.GenerateAdvertisingMaterials();
_publishingService.SubmitAdvertisement(@event.CampaignId,
advertisingMaterials);
}
public void Process(PublishingConfirmed @event)
{
var campaign = _repository.Load(@event.CampaignId);
campaign.TrackPublishingConfirmation(@event.ConfirmationId);
_repository.CommitChanges(campaign);
}
public void Process(PublishingRejected @event)
{
var campaign = _repository.Load(@event.CampaignId);
campaign.TrackPublishingRejection(@event.RejectionReason);
_repository.CommitChanges(campaign);
}
}
این Saga ساده State ندارد و به Messaging Infrastructure برای Delivery Event و اجرای Command متکی است. Sagaهای Stateful نیز وجود دارند؛ برای مثال برای Track کردن عملیات اجراشده تا در Failure بتوان Compensation درست را صادر کرد. در چنین حالتی Saga میتواند Event-Sourced Aggregate باشد و تاریخچهٔ کامل Eventهای دریافتی و Commandهای صادرشده را Persist کند. اجرای واقعی Command بهتر است از خود Saga بیرون و Asynchronous، شبیه Outbox، انجام شود:
public class CampaignPublishingSaga
{
private readonly ICampaignRepository _repository;
private readonly IList<IDomainEvent> _events;
...
public void Process(CampaignActivated activated)
{
var campaign = _repository.Load(activated.CampaignId);
var advertisingMaterials = campaign.GenerateAdvertisingMaterials();
var commandIssuedEvent = new CommandIssuedEvent(
target: Target.PublishingService,
command: new SubmitAdvertisementCommand(activated.CampaignId,
advertisingMaterials));
_events.Append(activated);
_events.Append(commandIssuedEvent);
}
public void Process(PublishingConfirmed confirmed)
{
var commandIssuedEvent = new CommandIssuedEvent(
target: Target.CampaignAggregate,
command: new TrackConfirmation(confirmed.CampaignId,
confirmed.ConfirmationId));
_events.Append(confirmed);
_events.Append(commandIssuedEvent);
}
public void Process(PublishingRejected rejected)
{
var commandIssuedEvent = new CommandIssuedEvent(
target: Target.CampaignAggregate,
command: new TrackRejection(rejected.CampaignId,
rejected.RejectionReason));
_events.Append(rejected);
_events.Append(commandIssuedEvent);
}
}
Outbox Relay برای هر CommandIssuedEvent، Command را روی Endpoint مربوط اجرا میکند. جداکردن تغییر State Saga از اجرای Command تضمین میکند Command حتی در صورت Crash فرایند بهطور قابل اتکا اجرا شود.
سازگاری
گرچه Saga تراکنشی چندجزئی را هماهنگ میکند، State اجزای درگیر Eventually Consistent است و هیچ دو Transactionای اتمیک محسوب نمیشوند. این موضوع با اصل Aggregate همراستاست:
فقط دادهٔ داخل مرز Aggregate را میتوان Strongly Consistent دانست؛ هر چیزی بیرون از مرز Eventually Consistent است.
این اصل را معیار قرار دهید تا از Saga برای جبران مرز Aggregate اشتباه سوءاستفاده نکنید. عملیات کسبوکاری که واقعاً به Strong Consistency مشترک نیاز دارند باید در همان Aggregate باشند.
Process Manager
Saga جریان ساده و خطی را مدیریت میکند؛ در تعریف دقیق، Event را به Command متناظر Match میکند. در مثالها:
- CampaignActivated → PublishingService.SubmitAdvertisement
- PublishingConfirmed → Campaign.TrackConfirmation
- PublishingRejected → Campaign.TrackRejection
Process Manager برای Business Processهایی است که تصمیمگیری و State دارند. بهعنوان یک واحد پردازش مرکزی، State توالی را نگه میدارد و مرحلهٔ بعد را تعیین میکند؛ شکل ۹-۱۴.
شکل ۹-۱۴ — Process Manager.
قاعدهٔ ساده: اگر Saga دارای if-else برای انتخاب مسیر بعدی باشد، احتمالاً در واقع Process Manager است.
تفاوت دیگر این است که Saga با مشاهدهٔ یک Event خاص بهطور ضمنی آغاز میشود؛ Process Manager معمولاً به یک Source Event واحد وابسته نیست و Business Process منسجمی با چند مرحله است، پس باید صریحاً Instantiate شود.
مثال سفر کاری: الگوریتم Routing ارزانترین مسیر پرواز را انتخاب و تأیید کارمند را میخواهد. اگر کارمند مسیر دیگری بخواهد، مدیر مستقیم باید تأیید کند. پس از رزرو پرواز، یکی از Hotelهای پیشتأییدشده برای تاریخ مناسب رزرو میشود. اگر Hotel موجود نباشد، بلیت پرواز باید Cancel شود. موجودیت مرکزیای که خودبهخود این فرایند را Trigger کند وجود ندارد؛ خود «رزرو سفر» فرایند است و باید با Process Manager پیاده شود؛ شکل ۹-۱۵.
شکل ۹-۱۵ — Process Manager رزرو سفر.
از دید پیادهسازی، Process Manager اغلب یک Aggregate مبتنی بر State یا Event-Sourced است:
public class BookingProcessManager
{
private readonly IList<IDomainEvent> _events;
private BookingId _id;
private Destination _destination;
private TripDefinition _parameters;
private EmployeeId _traveler;
private Route _route;
private IList<Route> _rejectedRoutes;
private IRoutingService _routing;
...
public void Initialize(Destination destination,
TripDefinition parameters,
EmployeeId traveler)
{
_destination = destination;
_parameters = parameters;
_traveler = traveler;
_route = _routing.Calculate(destination, parameters);
var routeGenerated = new RouteGeneratedEvent(
BookingId: _id,
Route: _route);
var commandIssuedEvent = new CommandIssuedEvent(
command: new RequestEmployeeApproval(_traveler, _route)
);
_events.Append(routeGenerated);
_events.Append(commandIssuedEvent);
}
public void Process(RouteConfirmed confirmed)
{
var commandIssuedEvent = new CommandIssuedEvent(
command: new BookFlights(_route, _parameters)
);
_events.Append(confirmed);
_events.Append(commandIssuedEvent);
}
public void Process(RouteRejected rejected)
{
var commandIssuedEvent = new CommandIssuedEvent(
command: new RequestRerouting(_traveler, _route)
);
_events.Append(rejected);
_events.Append(commandIssuedEvent);
}
public void Process(ReroutingConfirmed confirmed)
{
_rejectedRoutes.Append(route);
_route = _routing.CalculateAltRoute(destination,
parameters, rejectedRoutes);
var routeGenerated = new RouteGeneratedEvent(
BookingId: _id,
Route: _route);
var commandIssuedEvent = new CommandIssuedEvent(
command: new RequestEmployeeApproval(_traveler, _route)
);
_events.Append(confirmed);
_events.Append(routeGenerated);
_events.Append(commandIssuedEvent);
}
public void Process(FlightBooked booked)
{
var commandIssuedEvent = new CommandIssuedEvent(
command: new BookHotel(_destination, _parameters)
);
_events.Append(booked);
_events.Append(commandIssuedEvent);
}
...
}
در این مثال Process Manager شناسهٔ صریح و State Persistشوندهای دارد که سفر مورد رزرو را توصیف میکند. به Eventهای کنترلکنندهٔ Workflow مثل RouteConfirmed، RouteRejected و ReroutingConfirmed Subscribe میکند و CommandIssuedEvent تولید میکند تا Outbox Relay Command واقعی را اجرا کند.