Cantella-Worker-Role-Beanstalk

 view release on metacpan or  search on metacpan

lib/Cantella/Worker/Role/Beanstalk.pm  view on Meta::CPAN

has delete_on_max_tries => (
  is => 'ro',
  isa => 'Bool',
  required => 1,
  default => sub{ 0 },
);

after _start => sub {
  my $self = shift;
  for my $client ( @{ $self->beanstalk_clients } ){
    $client->disconnect;
    unless( $client->connect ){
      $self->error( $client->error );
    }
  }
};

sub get_work {
  my($self) = @_;
  for my $client ( shuffle @{ $self->beanstalk_clients }){
    if( my $job = $client->reserve($self->has_reserve_timeout ? $self->reserve_timeout : ())){
      if( $self->has_max_tries ){
        my $stats = $job->stats;
        if( $stats->reserves > $self->max_tries ){
          my $job_id = $job->id;
          my $tube = $stats->tube;
          my $args = join( ', ', map { "'$_'" } $job->args);
          if( $self->delete_on_max_tries ){
            $self->logger->notice("Job exceeds max-tries. Deleting job ${job_id} from tube '$tube' with args: $args");
            $job->delete;
          } else {
            $self->logger->notice("Job exceeds max-tries. Burying job ${job_id} from tube '$tube' with args: $args");
            $job->bury;
          }
          redo;
        }
      }
      return $job;
    } else {
      my $error = $client->error;
      if( $error ne 'TIMED_OUT'){
        $self->logger->error( $error );
      }
    }
  }
  return;
}

1;

__END__;

=head1 NAME

Cantella::Worker::Role::Beanstalk - Fetch Cantella::Worker jobs from beanstalkd

=head1 SYNOPSIS

    package TestWorkerPool;

    use Try::Tiny;
    use Moose;
    with(
      'Cantella::Worker::Role::Worker',
      'Cantella::Worker::Role::Beanstalk'
    );

    sub work {
      my ($self, $job) = @_;
      my @args = $job->args;
      try {
        if( do_something(@args) ){
          $job->delete; #work done successfully
        } else {
          $job->release({delay => 10}); #let's try again in 10 seconds
        }
      } catch {
        $job->bury; #job failed, bury it and log to file
        $self->logger->error("Burying job ".$job->id." due to error: '$_'");
      };
    }


=head1 ATTRIBUTES

=head2 beanstalk_clients

=over 4

=item B<beanstalk_clients> - reader

=back

Read-only, required, ArrayRef of L<Beanstalk::Client> instances.

=head2 reserve_timeout

=over 4

=item B<reserve_timeout> - reader

=item B<has_reserve_timeout> - predicate

=back

Read-only integer. The reserve timeout will be passed on to
L<Beanstalk::Client>'s C<reserve> method, and signals how long, in seconds,
the client should wait for a job to become available before timing out and
trying the next client in the pool.

B<WARNING:> If you only have one Beanstalk server, you might be tempted
to set not time out. B<Don't do this.> By setting no timeout, the reserve
command will block all other events, including signal handlers. Instead, it
is suggested that the C<reserve_timeout> is set to something that is resonable
for you workload and the load of your B<beanstalkd> process.

=head2 max_tries

=over 4

=item B<max_tries> - reader



( run in 1.873 second using v1.01-cache-2.11-cpan-5e09290becf )