Art*_*ess 7 c# windows-services scalability horizontal-scaling
我正在寻找有关如何扩展当前在我公司运行的Windows服务的一些输入.我们正在使用.NET 4.0(可以并且将来会在某个时候升级到4.5)并在Windows Server 2012上运行它.
关于服务
该服务的工作是查询日志表中的新行(我们正在使用Oracle数据库),处理信息,创建和/或更新其他5个表中的一堆行(让我们称之为跟踪表) ),更新日志表并重复.
日志记录表具有大量XML(每行最多可达20 MB),需要在其他5个跟踪表中进行选择和保存.始终以每小时500,000行的最大速率添加新行.
跟踪表的流量要高得多,从最小的一行中的90,000个新行到每小时最大表中可能有数百万行.更不用说那些表也有更新操作.
关于正在处理的数据
我觉得这一点对于根据这些对象的分组和处理方式找到解决方案非常重要.数据结构如下所示:
public class Report
{
public long Id { get; set; }
public DateTime CreateTime { get; set; }
public Guid MessageId { get; set; }
public string XmlData { get; set; }
}
public class Message
{
public Guid Id { get; set; }
}
Run Code Online (Sandbox Code Playgroud)
今天Windows服务我们几乎没有管理16核服务器上的负载(我不记得完整的规格,但可以肯定地说这台机器是野兽).我的任务是寻找扩展和添加更多机器的方法,这些机器将处理所有这些数据而不会干扰其他实例.
目前,每个Message都有自己的Thread并处理相关报告.我们批量处理报告,按其MessageId分组,以便在处理数据时将数据库查询的数量减少到最少.
限制
我正在寻找有关如何构建此类项目的任何意见或建议.我假设服务需要是无状态的,还是有办法以某种方式同步所有实例的缓存?我应该如何在所有实例之间进行协调,并确保它们不处理相同的数据?如何在它们之间平均分配负载?当然,如何处理实例崩溃而不完成它的工作?
编辑
删除了无关的信息
对于您的工作项,Windows Workflow可能是您重构服务的最快方法.
Windows Workflow Foundation @ MSDN
您将从WF中获得的最有用的事情是工作流持久性,如果工作流从保存它的最后一个点发生,则可以从持久点恢复正确设计的工作流.
这包括在处理工作流时,如果任何其他进程崩溃,则可以从另一个进程恢复工作流.如果使用共享工作流存储,则恢复过程不需要位于同一台计算机上.请注意,所有可恢复的工作流都需要使用工作流存储.
对于工作分配,您有几个选择.
通过WorkflowService类使用WCF端点通过工作流调用生成消息与基于主机的负载平衡相结合的服务.请注意,您可能希望在此处使用设计模式编辑器来构造入口方法,而不是手动设置Receive和相应的SendReply处理程序(这些映射到WCF方法).您可能会为每个消息调用该服务,也可能为每个Report调用该服务.请注意,该CanCreateInstance物业在这里很重要.与之绑定的每个调用都将创建一个独立运行的运行实例.
~
WorkflowService类(System.ServiceModel.Activities)@MSDN
接收类(System.ServiceModel.Activities)@ MSDN
Receive.CanCreateInstance属性(System.ServiceModel.Activities)@ MSDN
SendReply类(System.ServiceModel.Activities)@ MSDN
使用具有队列支持的服务总线.至少,您需要可能接受来自任意数量客户端的输入的内容,并且其输出可以唯一标识并仅处理一次.想到的一些是NServiceBus,MSMQ,RabbitMQ和ZeroMQ.在这里提到的项目中,NServiceBus完全是.NET开箱即用的.在云环境中,您的选项还包括特定于平台的产品,例如Azure Service Bus和Amazon SQS.
~
NServiceBus
MSMQ @ MSDN
RabbitMQ
ZeroMQ
Azure服务总线@ MSDN
亚马逊SQS @Amazon AWS
~
请注意,服务总线只是将启动消息的生产者和可以存在于任意数量的计算机上以从队列中读取的消费者之间的粘合剂.同样,您可以使用此间接进行报告生成.您的消费者将创建可以使用工作流持久性的工作流实例.
我通过自己编写所有这些可扩展性和冗余的东西来解决这个问题。如果有人需要的话,我将解释我做了什么以及我是如何做的。
我在每个实例中创建了一些进程来跟踪其他进程并了解特定实例可以处理哪些记录。启动时,实例将在数据库中(如果尚未)注册到名为 的表中Instances。该表具有以下列:
Id Number
MachineName Varchar2
LastActive Timestamp
IsMaster Number(1)
Run Code Online (Sandbox Code Playgroud)
MachineName在此表中注册并创建行后,如果未找到实例,实例将开始在单独的线程中每秒对该表执行 ping 操作,更新其LastActive列。然后,它从该表中选择所有行,并确保Master Instance(稍后详细介绍)仍然有效 - 这意味着它的LastActive时间在最后 10 秒内。如果主实例停止响应,它将接管控制权并将自己设置为主实例。在下一次迭代中,它将确保只有一个主实例(以防另一个实例决定同时承担控制权),如果不是,它将屈服于具有最低Id.
什么是主实例?
该服务的工作是扫描日志表并处理该数据,以便人们可以轻松地过滤和阅读它。我没有在我的问题中说明这一点,但它可能与这里相关。我们有一堆 ESB 服务器,每个请求都会将多条记录写入日志表,而我的服务的工作是近乎实时地跟踪它们。由于他们异步写入日志,我可能会在日志中获得finished processing request A之前的条目。started processing request A因此,我有一些代码可以对这些记录进行排序,并确保我的服务以正确的顺序处理数据。因为我需要扩展此服务,所以只有一个实例可以执行此逻辑,以避免大量不必要的数据库查询和可能的疯狂错误。
这就是 的用武之地Master Instance。只有它执行此排序逻辑并将日志记录 Id 临时保存在另一个名为 的表中ReportAssignment。该表的作用是跟踪哪些记录被处理以及由谁处理。处理完成后,记录将被删除。该表如下所示:
RecordId Number
InstanceId Number Nullable
Run Code Online (Sandbox Code Playgroud)
主实例对日志条目进行排序并在此处插入它们的 ID。我的所有服务实例都会以 1 秒的间隔检查此表,以查找未由任何人处理或正在由非活动实例处理的新记录以及([record's Id] % [number of isnstances] == [index of current instance in a sorted array of all the active instances]在 Pinging 过程中获取的)。该查询看起来有点像这样:
SELECT * FROM ReportAssignment
WHERE (InstanceId IS NULL OR InstanceId NOT IN (1, 2, 3)) // 1,2,3 are the active instances
AND RecordId % 3 == 0 // 0 is the index of the current instance in the list of active instances
Run Code Online (Sandbox Code Playgroud)
为什么我需要这样做?
RecordId % 3 == 1和RecordId % 3 == 2。 RecordId % [instanceCount] == [indexOfCurrentInstance]确保记录在所有实例之间均匀分布。 InstanceId NOT IN (1,2,3)允许实例接管崩溃实例正在处理的记录,并且在添加新实例时不处理已活动实例的记录。一旦实例查询这些记录,它将执行更新命令,将其设置InstanceId为自己的,并在日志记录表中查询具有这些 Id 的记录。处理完成后,它会删除 中的记录ReportAssignment。
总的来说我对此非常满意。它可以很好地扩展,确保在实例宕机时不会丢失数据,并且我们现有的代码几乎没有任何改变。