Showing posts with label Perl. Show all posts
Showing posts with label Perl. Show all posts

2013-04-23

Rolling your own message/job queue in MySQL (Part 3: Message Queue Implementation)

Note: These blog posts are in a stream-of-counciousness style with limited revision so I can rapidly progress without worrying about polishing the post.  If you notice a mistake, something missing, or even just confusing portions, let me know and I'll attempt to revise that portion.

The specification

Okay, in part 2 we created a rough specification for what we need to implement a message queue in mysql.  I'll summarize them here:

  1. A unique id for each message.
  2. A field for holding a unique transaction id (initially null) to prevent multiple dequeuing clients from colliding over the same message.
  3. One or more payload fields.
  4. Queuing a message is a row insertion.
  5. Dequeuing a message is a matter of an UPDATE with a LIMIT setting a unique transaction identifier on entries that do not have one, and then a subsequent SELECT for rows matching that unique transaction identifier.
  6. Accepting a message is a matter of deleting that row from the table.  If we need a history of messages, we will move the row to an archival table.
  7. Rejecting a message if achieved by resetting that message's unique transaction identifier back to null.
That's fairly straight-forward, and quite easy to implement.  For our implementation we'll use MySQL's AUTO_INCREMENT capability for the primary key to get a unique id for each message, and contrary to the example in part 2, we'll use a simple integer for the unique transaction id, as that should be sufficient for our needs.

The implementation

With the specification taken care of, our schema is extremely simple:

CREATE TABLE `queue_test` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT,
  `transaction_unique` mediumint(8) unsigned DEFAULT NULL,
  `payload` text,
  PRIMARY KEY (`id`),
  KEY `transaction_unique` (`transaction_unique`)
) ENGINE=InnoDB;

Now, it's perfectly valid to write SQL functions for queue and dequeue actions, and the accept and reject actions, but I'm not a fan of moving too much of my model code to the database.

While it does enforce an extra level of consistency on the data, I feel constrained architecturally and feel it won't scale as easily, since the database is generally much harder to scale than front-end servers in my experience.  Truthfully, there's probably a good trade-off point where there's real benefits from reducing round trips to the database by coding small quick routines as SQL functions, and our Dequeue implementation probably falls squarely into this.  I'm not going to implement that here though.

The code


Now, let's work up an implementation to this in an actual programming language.  I'm going to use Perl and DBIx::Class, because that's what I'm comfortable in and I believe it's a fairly compact but readable syntax.  If you have problems with that, feel free to translate into your favorite language.  Actually, please do, as that's the best way to get a feel for this.  For my following examples, please assume they all reside in the same file, even though I'll be presenting them in discrete chunks.

use 5.014; # implies strict
use warnings;
use utf8;

This is a fairly regular preamble in modern Perl.  We want to use newer features from version 5.14 (which implies "use strict;"), turn on warnings, and tell perl to expect possible UTF-8 data within the script (a sane default).

package Queue::Message {
    use base 'DBIx::Class::Core';
    __PACKAGE__->table( 'queue_test' );
    __PACKAGE__->add_columns(qw( id transaction_unique payload ));
    __PACKAGE__->set_primary_key( 'id' );
    __PACKAGE__->resultset_class( 'Queue::MessageBroker' );

    # To "accept" a message, we delete that row. This leaves the ORM object alone
    sub message_accept {
        my $self = shift;
        return $self->delete;
    }
    # To "reject" a message, we mark it as no longer part of a transaction
    sub message_reject {
        my $self = shift;
        return $self->update({ transaction_unique => undef })
    }

    # Convenience method to automatically convert the JSON payload
    sub message {
        my $self = shift;
        return decode_json( $self->payload );
    }

    1; # Class returns true
}

Here we have our Message class, which both defines our metadata for our table specification in DBIx::Class (in a shorthand syntax), and defines our message accept and reject methods, as defined in our specification.  Additionally, it provides a convenience method to decode the JSON we are serializing our message to so it can be easily stored.

package Queue::MessageBroker {
    use base 'DBIx::Class::ResultSet';
    use JSON;
    __PACKAGE__->load_components(
        qw(Helper::ResultSet Helper::ResultSet::Shortcut)
    );

    sub message_queue {
        my ($self,$message) = @_;
        return $self->create({ payload => encode_json( $message ) });
    }
    sub message_dequeue {
        my $self        = shift;
        my $wanted_msgs = shift || 1;
        my $uniq        = int(rand(2**24));
        # Mark some messages as part of this transaction
        my $result = $self->order_by('id')
                          ->rows($wanted_msgs)
                          ->search({ transaction_unique => undef })
                          ->update({ transaction_unique => $uniq });
        # Return marked messages
        return $self->search({ transaction_unique => $uniq })->all;
    }
    sub message_count { shift->search({ transaction_unique => undef })->count };

    1; # Class returns true
}

This is the definition of our Broker class, which as a DBIx::Class::ResultSet subclass handles the actual query operations.  Beyond some simple DBIx::Class setup for our inherited methods, plus a few convenience methods (chained order_by and rows methods), our Broker is responsible primarily for actually inserting (queuing), retrieving and deleting (dequeuing) messages, and as such methods to do those operations take up a majority of the implementation. 

package Queue {
    use base 'DBIx::Class::Schema';
    __PACKAGE__->load_classes( 'Message' );
    1; # Class returns true
}

Finally, we have our base ORM class.  This comes last because it expects to parse the object schema information from the Queue::Message class we defined above, and needs to know about any specifics we set in the ResultSet (MessageBroker) class.  Normally it can automatically find and load everything by itself from the separate files they are defined it, but since we are using a single file we have to be careful about the order we define classes.

That's it.  That the entire definition of our initial, somewhat naive implementation of a message queue on top of MySQL.  The only thing left to add is some code to drive the implementation and test that it works, so we might as well add that.

package main;
use JSON;
use Time::HiRes qw(time);

# Connect to DB with DBIx::Class schema object
my $schema = Queue->connect(
        "DBI:mysql:database=test;host=localhost;mysql_socket=/var/lib/mysql/mysql.sock",
        'queuetest_user', 'queuetest_pass', { PrintError => 0, RaiseError => 1, AutoCommit => 1 },
);
# Get our broker (DBIx::ResultSet for Message class)
my $broker = $schema->resultset('Message');

# Queue 4 messages with a high-resolution timestamp for data
$broker->message_queue({ time => time() }) for 1..4;
my @messages;

say 'Get one message, 4 -> 3';
say $broker->message_count, ' messages available';
@messages = $broker->message_dequeue;
say scalar(@messages) . " received";
printf("message: %d, time: %s\n", $_->id, $_->message->{time}) for @messages;
$_->message_accept for @messages;
print "\n";

say 'Get two messages, but reject them, 3 -> 3';
say $broker->message_count, ' messages available';
@messages = $broker->message_dequeue(2);
say scalar(@messages) . " received";
printf("message: %d, time: %s\n", $_->id, $_->message->{time}) for @messages;
$_->message_reject for @messages;
print "\n";

say 'Get two messages, 3 -> 1';
say $broker->message_count, ' messages available';
@messages = $broker->message_dequeue(2);
say scalar(@messages) . " received";
printf("message: %d, time: %s\n", $_->id, $_->message->{time}) for @messages;
$_->message_accept for @messages;
print "\n";

say 'Try to get two messages when 1 available, 1 -> 0';
say $broker->message_count, ' messages available';
@messages = $broker->message_dequeue(2);
say scalar(@messages) . " received";
printf("message: %d, time: %s\n", $_->id, $_->message->{time}) for @messages;
$_->message_accept for @messages;
print "\n";

Next time we'll look at performance, and see how quickly we can get messages in and out of this system.

2011-05-26

Distributed Workers

I have a project coming up where I'll need to utilize distributed workers. It's a bit odd in that the workers will come and go, and they'll most likely need a copy of the subset of the data they are working on so they can efficiently process it, but the workload is very time dependent. Put another way, I need to keep data synchronized in some fashion between the master server and the remote client in a way that is testable.

I'm thinking I'm going to see problems where two works are working on the same data set but in doing so get slightly different results. Returning different results is not only possible, but probable, since processing the data requires them check a resource that sometimes flaps between values when in transition, for minutes at a time.

Here's the criteria for the system as I see it so far:

Server:

  • Canonical data source; Data stored in some sort of DB

  • Accepts registrations from clients/workers

  • Creates jobs/tasks in a work queue

  • Assigns jobs/tasks from work queue to registered workers

  • Accepts results from workers or times out task after appropriate wait



Client (worker):

  • Mostly shared code base (re-use modules defining data as objects)

  • Registers with server

  • Accepts tasks from server, processes data, returns result

  • Keeps copy of current set of data it is responsible for processing, only returns changes to data, not whole update



Here's what I'm wondering:

How much of this is based in my assumptions for what I'll need underneath? I've already thought of the DB structure needed to support this, and how I'll link between all the structures in the data. If I assume I'm using some sort of NoSQL solution, such as MongoDB, CouchDB (or whatever it's called now), or something else, are there assumptions I can make about the system that reduces complexity?

Are there modules available (preferably in Perl) to manage some of the work assignment tasks for me?

I would prefer to pass object state back and forth for the tasks. I can imagine passing an object name and a way to initialize that object to the state defined, that's not too hard. I DO want to have the objects that I'm passing easily abstracted to the DB on the server side. If I have the workers contain the same object code, can I do that without requiring the client deal with DB code? That is, can I easily abstract the object ORM layer out from the client? Maybe with roles using Moose?

If I use Moose, I know there's a startup speed penalty, which is not a problem. I'm more worried about any execution inefficiencies, since this is time dependent (to a sub-second level, but not quite ms dependent level. I haven't had a chance to use Moose in a project yet, so I'm not aware of the specifics. I do hear it's tunable so I can omit features for speed, which is a nice trade-off.

Some representation for the changes in a data structure, or just JSON if that works as a common format, would be very useful. If I can find a module that provides this, great. Otherwise, I suspect I'll be writing my own after researching data diffs.

In any case, I'll update here as I come to conclusions or find solutions.

2011-05-09

Perlbrew to the rescue

Chromatic recently posted about the support lifetime of Perl, and it's extension through enterprise distributions. While I don't particularly buy his arguments against enterprise need for back-patching and supporting older versions of Perl (and I suspect neither does he completely. He always strikes be as somewhat of a provocateur, a noble profession), I do agree that App::perlbrew is part of the solution.

While we seem to be in agreement that perlbrew is the solution, he seems to think (it's ambiguous in the post) that perlbrew can't be included in existing enterprise releases, such as RHEL 5 (which I'm most familiar with, and will restrict my examples to). I don't see any major reason that it can't be made available to existing enterprise distributions (in a supported manner, even). RHEL has a long history of providing feature enhancements and new packages/programs in their point releases (as opposed to the strictly bug and security fixes between point releases), so including perlbrew would be easily accomplished. Even if RHEL doesn't want to include it for whatever reason, getting it included in CentOS through their extras would be trivial (well, as trivial as doing anything with the CentOS developers is these days), and provide a real enhancement.

Of course, the Perl versions installed from perlbrew themselves would not necessarily be supported, but that's an easy point of demarcation to define. Different support policies could (and should) be defined for applications developed and/or deployed on a platform, as opposed to the platform itself. This allows for easy updating of subsystems that aren't related to the deployed application, but are required for security reasons.

I remember the most recent time I was migrating (and updating) RT. I spent quite a while swimming through dependency hell, making RPMs of all the required CPAN modules that weren't already available through RPMForge, EPEL and the like. During the final stages of that, visions of making my own RT bundle that auto installed Perl through perlbrew, along with the latest relevant CPAN modules were definitely dancing through my head. I think the world is a more barren place for my lack of motivation after the migration project was complete.

Now, for anyone who really just doesn't get why enterprise distributions need to keep the old version of Perl around, consider the following; enterprise distributions need the ability to ensure that during any update, nothing can or will go wrong (at least as much as they can). This often means limiting an installed program to the original shipped version, and back-porting non-conflicting features. When you have to upgrade hundreds of servers, this is essential. This is the continual struggle between system administrators and developers. Both groups are striving for stability, maintainability, and security, but these concepts mean slightly different things to each group. The beauty of perlbrew is it allows each group to have their own sandbox that they can correctly apply their goals to.

Please note than while I'm not sure the support requirements of perlbrew itself, if they aren't as minimal as possible to run on as old a version of Perl as possible, than I see that as a serious design flaw. I don't believe this to apply to CPAN modules in general though.

2010-07-21

Threading: Perl vs Python


Apparently some people don't believe that python's threading, while easier, is necessarily slower than a system that uses OS threads. This illustrates a fundamental misunderstanding of how OS threads work compared to "green" threads, also called coroutines. Python's threading model doesn't seem to quite match green threads exactly, but because of the global interpreter lock, they have a lot of the same performance problems, so I'll treat this test as an indication of the coroutine threading model's performance as well.

The problem with green threads is that they don't scale with the hardware. Sure, they scale beautifully in the number of processes you can spawn and control, but just because it shows two concurrent processes, doesn't mean you are actually utilizing the hardware effectively. Python's global interpreter lock ensures that two python processes can't be running at the same time. This resolves any possible variable locking problems. It also means, only one CPU core is being utilized, which has clear performance implications.

To illustrate this, I've created two simple test programs, one in python, one in Perl. They each split into a specifiable number of threads, and each thread performs some repetitive math operations 10,000,000 times and reports back how long it took to a centralized data structure. Finally, it prints out the sum of the times the threads took to perform the task, and the time the main process actually took spawn and collect them, which is the real processing time. Here they are:

Perl:
use strict;
use warnings;
use threads;
use threads::shared;
use Time::HiRes qw/time/;
my @thread_times :shared;

sub workfunc {
my $start = time();
print "Thread ".threads::tid()." :: started at ".localtime()."\n";
my $init = shift;
for (1 .. 10_000_000) {
$init += (($_%4)*($_%4)) / ($_%10+1);
}
my $runtime = time() - $start;
{ lock(@thread_times); push(@thread_times, $runtime); }
print "Thread ".threads::tid()." :: stopped at ".localtime()."\n";
return $runtime;
}

# Setup
my $num_threads = shift || 2;
print "Running with $num_threads threads\n";

# Main code
my $master_start = time();
my @threads = map { threads->create(\&workfunc, rand(10)); } 1 .. $num_threads;
my $thread_time = 0;
for my $t (@threads) {
$thread_time += $t->join();
}
my $master_time = time() - $master_start;

# Find and print times
my $lock_time = 0;
{ lock(@thread_times); $lock_time += $_ for @thread_times; }
print "Master logged $master_time runtime\n";
print "Threads reported $lock_time combined runtime\n";


Python:
import math
import random
import time
import sys
import threading
thread_times = []

def workfunc(id,init):
l = threading.local()
l.start = time.clock()
print "Thread "+str(id)+" :: started at "+time.asctime()
for i in xrange (1, 10000000):
init += ((i%4)*(i%4)) / (i%10+1)
l.runtime = time.clock() - l.start
thread_times.append(l.runtime)
print "Thread "+str(id)+" :: stopped at "+time.asctime()
return l.runtime

# Setup
num_threads = 2
if len(sys.argv)>1:
num_threads = int(sys.argv[1])
print "Running with "+str(num_threads)+" threads"

# Main code
master_start = time.clock()
threads = []
for i in xrange (1, num_threads+1):
t = threading.Thread(target=workfunc,args=(i,random.random()*10))
t.start()
threads.append(t)
for t in threads:
t.join()
master_time = time.clock() - master_start

# Find and print times
thread_time = 0
for t in thread_times:
thread_time += t
print "Master logged "+str(master_time)+" runtime";
print "Threads reported "+str(thread_time)+" combined runtime";
The results are very telling, and here they are:

Perl:
$ perl thread.pl 1
Running with 1 threads
Thread 1 :: started at Wed Jul 21 12:47:26 2010
Thread 1 :: stopped at Wed Jul 21 12:47:31 2010
Master logged 5.21124505996704 runtime
Threads reported 5.20653009414673 combined runtime
$ perl thread.pl 2
Running with 2 threads
Thread 1 :: started at Wed Jul 21 12:47:36 2010
Thread 2 :: started at Wed Jul 21 12:47:36 2010
Thread 2 :: stopped at Wed Jul 21 12:47:42 2010
Thread 1 :: stopped at Wed Jul 21 12:47:42 2010
Master logged 5.42760705947876 runtime
Threads reported 10.6412711143494 combined runtime
$ perl thread.pl 3
Running with 3 threads
Thread 1 :: started at Wed Jul 21 12:47:47 2010
Thread 2 :: started at Wed Jul 21 12:47:47 2010
Thread 3 :: started at Wed Jul 21 12:47:47 2010
Thread 2 :: stopped at Wed Jul 21 12:47:54 2010
Thread 1 :: stopped at Wed Jul 21 12:47:54 2010
Thread 3 :: stopped at Wed Jul 21 12:47:56 2010
Master logged 8.67003512382507 runtime
Threads reported 22.7089991569519 combined runtime
$ perl thread.pl 4
Running with 4 threads
Thread 1 :: started at Wed Jul 21 12:47:59 2010
Thread 2 :: started at Wed Jul 21 12:47:59 2010
Thread 3 :: started at Wed Jul 21 12:47:59 2010
Thread 4 :: started at Wed Jul 21 12:47:59 2010
Thread 1 :: stopped at Wed Jul 21 12:48:09 2010
Thread 3 :: stopped at Wed Jul 21 12:48:09 2010
Thread 2 :: stopped at Wed Jul 21 12:48:10 2010
Thread 4 :: stopped at Wed Jul 21 12:48:10 2010
Master logged 10.9765570163727 runtime
Threads reported 43.3241600990295 combined runtime
Note that a single thread takes approximately 5.2 seconds. Two threads takes only a few fractions of a second more. That's because each thread ran on it's own core. The threads report they had a combined run time of 10+ seconds, as each had a 5+ second run time, but it was *in parallel*, so the real processing time was ~5.2 seconds. My system is a dual core system, so we finally see some slowdown at three threads, where it takes 50% longer. Four threads takes twice as long as two threads, as expected.

Now let's look at Python:
$ python thread.py 1
Running with 1 threads
Thread 1 :: started at Wed Jul 21 12:53:03 2010
Thread 1 :: stopped at Wed Jul 21 12:53:08 2010
Master logged 5.25 runtime
Threads reported 5.25 combined runtime
$ python thread.py 2
Running with 2 threads
Thread 1 :: started at Wed Jul 21 12:53:10 2010
Thread 2 :: started at Wed Jul 21 12:53:10 2010
Thread 2 :: stopped at Wed Jul 21 12:53:22 2010
Thread 1 :: stopped at Wed Jul 21 12:53:22 2010
Master logged 14.66 runtime
Threads reported 29.09 combined runtime
$ python thread.py 3
Running with 3 threads
Thread 1 :: started at Wed Jul 21 12:53:25 2010
Thread 2 :: started at Wed Jul 21 12:53:25 2010
Thread 3 :: started at Wed Jul 21 12:53:25 2010
Thread 3 :: stopped at Wed Jul 21 12:53:43 2010
Thread 1 :: stopped at Wed Jul 21 12:53:44 2010
Thread 2 :: stopped at Wed Jul 21 12:53:44 2010
Master logged 22.48 runtime
Threads reported 66.26 combined runtime
$ python thread.py 4
Running with 4 threads
Thread 1 :: started at Wed Jul 21 12:53:46 2010
Thread 2 :: started at Wed Jul 21 12:53:46 2010
Thread 3 :: started at Wed Jul 21 12:53:46 2010
Thread 4 :: started at Wed Jul 21 12:53:46 2010
Thread 2 :: stopped at Wed Jul 21 12:54:12 2010
Thread 1 :: stopped at Wed Jul 21 12:54:12 2010
Thread 3 :: stopped at Wed Jul 21 12:54:13 2010
Thread 4 :: stopped at Wed Jul 21 12:54:13 2010
Master logged 30.16 runtime
Threads reported 120.4 combined runtime
Here we see where the global interpreter lock causing us problems. Each additional thread causes it to take the same amount more time. That's because only a single thread can run at any time, yielding no time savings for running four threads. It's equivalent to running them in sequence, in this case.

Now, it's only fair to note, this isn't always the case. We are seeing this behavior because this is a CPU bound problem. If the workload for each thread were IO bound, we would probably see much better performance from python than we are in this case. Green threads have the ability to much more closely and easily control the separate threads, which yields good results when you don't actually NEED multiple CPU cores.

In summary, Green threads, coroutines, and python's implementation of threads have clear disadvantages under certain workloads, to the point of having negligible to no performance benefits in some cases. If you are using a language that only offers one of these methods of threading and you have a process model that lends itself well to threading, I feel for you.

Update: I figured I should post the python and perl versions before I'm asked/accused.
$ perl -v
This is perl, v5.10.1 (*) built for x86_64-linux-thread-multi
$ python -V
Python 2.6.2

Both are the standard RPMs included with the RHEL6 beta 2 install.

$ rpm -q perl
perl-5.10.1-109.el6.x86_64
$ rpm -q python
python-2.6.2-7.el6.x86_64