MT#18663 MT#20893 row bulk processing framework WIP #7

+allow #* symbols in setoptionitems in
 Features_Define.cfg
+prevent double execution+log properly
+make sure to close all db connections before
 forking threads, not only project-specific ones
+per-db dao folder for ipgallery migration project
+start with ngcp db dao's:
 +contract_balances
 +contracts
+refactor dao/sql processing class
 +params hash for method with long arg list
 +take out obselete methods
 +dedicated unique field option for "insert_record"
 +split SqlRecord class into
  SqlRecord&SqlProcessor
+dependencies in Build.pl and debian/control
 +took out comments from debian/control
+calcualte deltas of Subscriber_Define.cfg imports
+"truncate subscriber import" task

Change-Id: I8eab88319e6a4ba3d5b6ba935e50a56cc1763ad9
changes/43/6943/8
Rene Krenn 10 years ago
parent 1c669a3a73
commit b4f244a9f1

@ -0,0 +1,74 @@
use Module::Build qw();
my $builder = Module::Build->new(
dist_name => 'NGCP::BulkProcessor',
dist_abstract => 'Framework for parallel/distributed processing of record blocks',
license => 'GPL_3',
dist_author => 'Rene Krenn <rkrenn@sipwise.com>',
dist_version_from => 'lib/NGCP/BulkProcessor/Globals.pm',
perl_version_from => 'lib/NGCP/BulkProcessor/Globals.pm',
requires => {
'Archive::Zip' => 0,
'OLE::Storage_Lite' => 0,
'XML::LibXML::Reader' => 0,
'Email::MIME' => 0,
'Email::MIME::Attachment::Stripper' => 0,
'URI::Find' => 0,
'LWP::UserAgent' => 0,
'HTTP::Request' => 0,
'DateTime' => 0,
'Time::HiRes' => 0,
'Time::Warp' => 0,
'DateTime::TimeZone' => 0,
'DateTime::Format::Strptime' => 0,
'DateTime::Format::ISO8601' => 0,
'Tie::IxHash' => 0,
'YAML::Tiny' => 0,
'Log::Log4Perl' => 0,
'MIME::Base64' => 0,
'MIME::Lite' => 0,
'Net::SMTP' => 0,
'JSON::XS' => 0,
'Data::Dump' => 0,
'YAML::XS' => 0,
'XML::Dumper' => '0.81',
'PHP::Serialization' => 0,
#'Gearman::Worker' => 0,
'Gearman::Client' => 0,
#'Gearman::Task' => 0,
'Digest::MD5' => 0,
'Data::UUID' => 0,
'Net::Address::IP::Local' => 0,
'Date::Manip' => 0,
'Date::Calc' => 0,
#Sys::CpuAffinity
'Marpa::R2' => 0,
'Data::Dumper::Concise' => 0,
'IO::Socket::SSL' => 0,
'Mail::IMAPClient' => 0,
'DBI' => '1.608',
'DBD::CSV' => '0.26',
'Locale::Recode' => 0,
'Spreadsheet::ParseExcel' => 0,
'Spreadsheet::ParseExcel::FmtUnicode' => 0,
'Text::CSV_XS' => 0,
'MIME::Parser' => 0,
'HTML::Entities' => 0,
'IO::Uncompress::Unzip' => 0,
'DBD::mysql' => '4.014',
#'DBD::Oracle' => '1.21',
#'DBD::Pg' => '2.17.2',
#'DBD::ODBC' => '1.50',
'DBD::SQLite' => '1.29',
},
test_requires => {
'Module::Runtime' => 0,
'Test::Unit::Procedural' => 0,
},
add_to_cleanup => ['NGCP-BulkProcessor-*', 'Excel-Reader-*'],
);
$builder->create_build_script;

@ -0,0 +1,40 @@
=encoding UTF-8
Bulk-Processor
=head1 NAME
README - basic information for users prior to downloading
=head1 INSTALLATION
See L<http://www.cpan.org/modules/INSTALL.html>.
=head1 DEPENDENCIES
See distribution meta file.
=head1 LICENCE
This software is Copyright © 2016 by Sipwise GmbH, Austria.
This program is free software; you can redistribute it
and/or modify it under the terms of the GNU General Public
License as published by the Free Software Foundation; either
version 3 of the License, or (at your option) any later
version.
This program is distributed in the hope that it will be
useful, but WITHOUT ANY WARRANTY; without even the implied
warranty of MERCHANTABILITY or FITNESS FOR A PARTICULAR
PURPOSE. See the GNU General Public License for more
details.
You should have received a copy of the GNU General Public
License along with this package; if not, write to the Free
Software Foundation, Inc., 51 Franklin St, Fifth Floor,
Boston, MA 02110-1301 USA
On Debian systems, the full text of the GNU General Public
License version 3 can be found in the file
F</usr/share/common-licenses/GPL-3>.

48
debian/control vendored

@ -11,6 +11,50 @@ Package: ngcp-bulk-processor-pro
Architecture: all
Depends:
perl,
libarchive-zip-perl,
libole-storage-lite-perl,
libxml-libxml-perl,
libemail-mime-perl,
libemail-mime-attachment-stripper-perl,
liburi-find-perl,
libwww-perl,
libdatetime-perl,
libtime-warp-perl,
libdatetime-timezone-perl,
libdatetime-format-strptime-perl,
libdatetime-format-iso8601-perl,
libtie-ixhash-perl,
libyaml-tiny-perl,
liblog-log4perl-perl,
libmime-base64-perl,
libmime-lite-perl,
libnet-smtp-ssl-perl,
libjson-xs-perl,
libdata-dump-perl,
libyaml-libyaml-perl,
libxml-dumper-perl,
libphp-serialization-perl,
libgearman-client-perl,
libdigest-md5-perl,
libdata-uuid-perl,
libnet-address-ip-local-perl,
libdate-manip-perl,
libdate-calc-perl,
libmarpa-r2-perl,
libdata-dumper-concise-perl,
libio-socket-ssl-perl,
libmail-imapclient-perl,
libdbd-csv-perl,
libintl-perl,
libspreadsheet-parseexcel-perl,
libtext-csv-xs-perl,
libmime-tools-perl,
libemail-mime-perl,
libhtml-parser-perl,
libio-compress-perl,
libdbd-mysql-perl,
libdbd-sqlite3-perl,
${misc:Depends},
Description: NGCP bulk-processor tool
Tool for bulk process for NGCP.
${perl:Depends}
Description: NGCP bulk processor framework
Framework for parallel/distributed processing of record blocks.

@ -37,7 +37,7 @@ use NGCP::BulkProcessor::SqlConnectors::CSVDB;
#use NGCP::BulkProcessor::SqlConnectors::SQLServerDB;
use NGCP::BulkProcessor::RestConnectors::NGCPRestApi;
use NGCP::BulkProcessor::SqlRecord qw(cleartableinfo);
use NGCP::BulkProcessor::SqlProcessor qw(cleartableinfo);
use NGCP::BulkProcessor::Utils qw(threadid);

@ -0,0 +1,110 @@
package NGCP::BulkProcessor::Dao::Trunk::billing::billing_mappings;
use strict;
## no critic
#use File::Basename;
#use Cwd;
#use lib Cwd::abs_path(File::Basename::dirname(__FILE__) . '/../../../');
use NGCP::BulkProcessor::ConnectorPool qw(
get_billing_db
);
use NGCP::BulkProcessor::SqlProcessor qw(
checktableinfo
insert_record
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
our @EXPORT_OK = qw(
gettablename
check_table
insert_row
);
#my $logger = getlogger(__PACKAGE__);
my $tablename = 'billing_mappings';
my $get_db = \&get_billing_db;
my $expected_fieldnames = [
'id',
'start_date',
'end_date',
'billing_profile_id',
'contract_id',
'product_id',
'network_id',
];
my $indexes = {};
# 'balance_interval' => [ 'contract_id','start','end' ],
# 'invoice_idx' => [ 'invoice_id' ],
#};
my $insert_unique_fields = []; #[ 'contract_id','start','end' ];
sub new {
my $class = shift;
my $self = NGCP::BulkProcessor::SqlRecord->new($get_db,
$tablename,
$expected_fieldnames,$indexes);
bless($self,$class);
copy_row($self,shift,$expected_fieldnames);
return $self;
}
sub insert_row {
my ($data,$insert_ignore) = @_;
check_table();
return insert_record($get_db,$tablename,$data,$insert_ignore,$unique_fields) = @_;
}
sub buildrecords_fromrows {
my ($rows,$load_recursive) = @_;
my @records = ();
my $record;
if (defined $rows and ref $rows eq 'ARRAY') {
foreach my $row (@$rows) {
$record = __PACKAGE__->new($row);
# transformations go here ...
push @records,$record;
}
}
return \@records;
}
sub gettablename {
return $tablename;
}
sub check_table {
return checktableinfo($get_db,
$tablename,
$expected_fieldnames,
$indexes);
}
1;

@ -9,27 +9,27 @@ use strict;
use NGCP::BulkProcessor::ConnectorPool qw(
get_billing_db
billing_db_tableidentifier
);
use NGCP::BulkProcessor::SqlRecord qw(checktableinfo);
use NGCP::BulkProcessor::SqlProcessor qw(
checktableinfo
insert_record
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
our @EXPORT_OK = qw(
XX
backoffice_client_byboclientid
sync_table
drop_table
check_local_table
check_source_table);
gettablename
check_table
insert_row
);
#my $logger = getlogger(__PACKAGE__);
my $tablename = 'contract_balances';
my $get_db = \&get_billing_db;
my $get_tablename = \&billing_db_tableidentifier;
my $expected_fieldnames = [
'id',
@ -47,14 +47,19 @@ my $expected_fieldnames = [
'underrun_lock',
];
#my $indexes = { $tablename . '_subscribernumber' => ['subscribernumber(11)'] };
my $indexes = {};
# 'balance_interval' => [ 'contract_id','start','end' ],
# 'invoice_idx' => [ 'invoice_id' ],
#};
my $insert_unique_fields = []; #[ 'contract_id','start','end' ];
sub new {
my $class = shift;
my $self = NGCP::BulkProcessor::SqlRecord->new($get_db,
gettablename(),
$expected_fieldnames);
$tablename,
$expected_fieldnames,$indexes);
bless($self,$class);
@ -64,6 +69,13 @@ sub new {
}
sub insert_row {
my ($data,$insert_ignore) = @_;
check_table();
return insert_record($get_db,$tablename,$data,$insert_ignore,$unique_fields) = @_;
}
sub buildrecords_fromrows {
@ -88,15 +100,16 @@ sub buildrecords_fromrows {
sub gettablename {
return &$get_tablename($get_db,$tablename);
return $tablename;
}
sub check_table {
return checktableinfo($get_db,
gettablename(),
$expected_fieldnames);
$tablename,
$expected_fieldnames,
$indexes);
}

@ -0,0 +1,122 @@
package NGCP::BulkProcessor::Dao::Trunk::billing::contracts;
use strict;
## no critic
#use File::Basename;
#use Cwd;
#use lib Cwd::abs_path(File::Basename::dirname(__FILE__) . '/../../../');
use NGCP::BulkProcessor::ConnectorPool qw(
get_billing_db
);
use NGCP::BulkProcessor::SqlProcessor qw(
checktableinfo
insert_record
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
our @EXPORT_OK = qw(
gettablename
check_table
insert_row
);
#my $logger = getlogger(__PACKAGE__);
my $tablename = 'contracts';
my $get_db = \&get_billing_db;
my $expected_fieldnames = [
'id',
'customer_id',
'contact_id',
'order_id',
'profile_package_id',
'status',
'external_id',
'modify_timestamp',
'create_timestamp',
'activate_timestamp',
'terminate_timestamp',
'max_subscribers',
'send_invoice',
'subscriber_email_template_id',
'passreset_email_template_id',
'invoice_email_template_id',
'invoice_template_id',
'vat_rate',
'add_vat',
];
my $indexes = {};
# 'balance_interval' => [ 'contract_id','start','end' ],
# 'invoice_idx' => [ 'invoice_id' ],
#};
my $insert_unique_fields = []; #[ 'contract_id','start','end' ];
sub new {
my $class = shift;
my $self = NGCP::BulkProcessor::SqlRecord->new($get_db,
$tablename,
$expected_fieldnames,$indexes);
bless($self,$class);
copy_row($self,shift,$expected_fieldnames);
return $self;
}
sub insert_row {
my ($data,$insert_ignore) = @_;
check_table();
return insert_record($get_db,$tablename,$data,$insert_ignore,$unique_fields) = @_;
}
sub buildrecords_fromrows {
my ($rows,$load_recursive) = @_;
my @records = ();
my $record;
if (defined $rows and ref $rows eq 'ARRAY') {
foreach my $row (@$rows) {
$record = __PACKAGE__->new($row);
# transformations go here ...
push @records,$record;
}
}
return \@records;
}
sub gettablename {
return $tablename;
}
sub check_table {
return checktableinfo($get_db,
$tablename,
$expected_fieldnames,
$indexes);
}
1;

@ -3,6 +3,7 @@ use strict;
## no critic
use DateTime qw();
use Time::HiRes qw(); #prevent warning from Time::Warp
use Time::Warp qw();
use DateTime::TimeZone qw();

@ -88,7 +88,19 @@ sub process {
my $self = shift;
my ($file,$process_code,$init_process_context_code,$uninit_process_context_code,$multithreading) = @_;
my %params = @_;
my ($file,
$process_code,
$init_process_context_code,
$uninit_process_context_code,
$multithreading) = @params{qw/
file
process_code
init_process_context_code
uninit_process_context_code
multithreading
/};
#my ($file,$process_code,$init_process_context_code,$uninit_process_context_code,$multithreading) = @_;
if (ref $process_code eq 'CODE') {

@ -107,7 +107,8 @@ umask 0000;
# general constants
our $system_name = 'Sipwise Bulk Processing Framework';
our $system_version = '0.0.1'; #keep this filename-save
our $VERSION = '0.0.1';
our $system_version = $VERSION; #keep this filename-save
our $system_abbreviation = 'sbpf'; #keep this filename-, dbname-save
our $system_instance = 'initial'; #'test'; #'2014'; #dbname-save 0-9a-z_
our $system_instance_label = 'test';

@ -509,9 +509,9 @@ sub rowinsertskipped {
sub rowupdateskipped {
my ($db,$tablename,$logger) = @_;
my ($db,$tablename,$matched,$logger) = @_;
if (defined $logger) {
$logger->info(_getsqlconnectorinstanceprefix($db) . 'row update skipped');
$logger->info(_getsqlconnectorinstanceprefix($db) . "row update skipped, $matched matching rows");
}
}

@ -1,4 +1,4 @@
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption;
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption;
use strict;
## no critic
@ -13,7 +13,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw
);
#import_db_tableidentifier
use NGCP::BulkProcessor::SqlRecord qw(
use NGCP::BulkProcessor::SqlProcessor qw(
registertableinfo
create_targettable
checktableinfo
@ -21,8 +21,9 @@ use NGCP::BulkProcessor::SqlRecord qw(
insert_stmt
);
use NGCP::BulkProcessor::SqlRecord qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
@ -41,14 +42,15 @@ my $get_db = \&get_import_db;
#my $get_tablename = \&import_db_tableidentifier;
my $expected_fieldnames = [ 'subscribernumber',
'option'];
my $expected_fieldnames = [
'subscribernumber',
'option'
];
# table creation:
my $primarykey_fieldnames = [ 'subscribernumber', 'option' ];
my $indexes = {};
my $fixtable_statements = [];
#my $fixtable_statements = [];
sub new {
@ -127,7 +129,7 @@ sub buildrecords_fromrows {
# transformations go here ...
if ($load_recursive) {
$record->{_optionsetitems} = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem::findby_subscribernumber_option(
$record->{_optionsetitems} = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::findby_subscribernumber_option(
$record->{subscribernumber},
$record->{option},
$load_recursive

@ -1,4 +1,4 @@
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem;
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem;
use strict;
## no critic
@ -13,7 +13,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw
);
#import_db_tableidentifier
use NGCP::BulkProcessor::SqlRecord qw(
use NGCP::BulkProcessor::SqlProcessor qw(
registertableinfo
create_targettable
checktableinfo
@ -21,6 +21,7 @@ use NGCP::BulkProcessor::SqlRecord qw(
insert_stmt
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
@ -39,15 +40,18 @@ my $get_db = \&get_import_db;
#my $get_tablename = \&import_db_tableidentifier;
my $expected_fieldnames = [ 'subscribernumber',
'option',
'optionsetitem' ];
my $expected_fieldnames = [
'subscribernumber',
'option',
'optionsetitem'
];
# table creation:
my $primarykey_fieldnames = []; #[ 'subscribernumber', 'option', 'optionsetitem' ];
my $indexes = { $tablename . '_subscribernumber_option_optionsetitem' => ['subscribernumber(11)', 'option(32)', 'optionsetitem(32)'] }; #(25),(27)
my $fixtable_statements = [];
my $indexes = {
$tablename . '_subscribernumber_option_optionsetitem' => ['subscribernumber(11)', 'option(32)', 'optionsetitem(32)'], #(25),(27)
};
#my $fixtable_statements = [];
sub new {

@ -1,4 +1,4 @@
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp;
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp;
use strict;
## no critic
@ -13,7 +13,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw
);
#import_db_tableidentifier
use NGCP::BulkProcessor::SqlRecord qw(
use NGCP::BulkProcessor::SqlProcessor qw(
registertableinfo
create_targettable
checktableinfo
@ -21,6 +21,7 @@ use NGCP::BulkProcessor::SqlRecord qw(
insert_stmt
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
@ -46,11 +47,10 @@ my $expected_fieldnames = [
'lrn_code',
];
# table creation:
my $primarykey_fieldnames = [ 'lrn_code', 'ported_number' ];
my $indexes = {};
my $fixtable_statements = [];
#my $fixtable_statements = [];
sub new {

@ -1,4 +1,4 @@
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber;
package NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber;
use strict;
## no critic
@ -13,7 +13,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw
);
#import_db_tableidentifier
use NGCP::BulkProcessor::SqlRecord qw(
use NGCP::BulkProcessor::SqlProcessor qw(
registertableinfo
create_targettable
checktableinfo
@ -21,6 +21,7 @@ use NGCP::BulkProcessor::SqlRecord qw(
insert_stmt
);
use NGCP::BulkProcessor::SqlRecord qw();
require Exporter;
our @ISA = qw(Exporter NGCP::BulkProcessor::SqlRecord);
@ -29,9 +30,17 @@ our @EXPORT_OK = qw(
gettablename
check_table
getinsertstatement
getupsertstatement
findby_subscribernumber
countby_subscribernumber
update_delta
findby_delta
countby_delta
$deleted_delta
$updated_delta
$added_delta
);
my $tablename = 'subscriber';
@ -50,13 +59,17 @@ my $expected_fieldnames = [
'time_zone_name', #malta
'lang_code', #eng
'barring_profile', #None
'delta',
];
# table creation:
my $primarykey_fieldnames = [ 'country_code', 'area_code', 'dial_number' ];
my $indexes = { $tablename . '_delta' => [ 'delta(7)' ]};
#my $fixtable_statements = [];
my $indexes = {};
my $fixtable_statements = [];
our $deleted_delta = 'DELETED';
our $updated_delta = 'UPDATED';
our $added_delta = 'ADDED';
sub new {
@ -84,6 +97,27 @@ sub create_table {
}
sub findby_delta {
my ($delta,$load_recursive) = @_;
check_table();
my $db = &$get_db();
my $table = $db->tableidentifier($tablename);
return [] unless defined $delta;
my $rows = $db->db_get_all_arrayref(
'SELECT * FROM ' .
$table .
' WHERE ' .
$db->columnidentifier('delta') . ' = ?'
,$delta);
return buildrecords_fromrows($rows,$load_recursive);
}
sub findby_subscribernumber {
my ($subscribernumber,$load_recursive) = @_;
@ -129,6 +163,49 @@ sub countby_subscribernumber {
}
sub countby_delta {
my ($delta) = @_;
check_table();
my $db = &$get_db();
my $table = $db->tableidentifier($tablename);
my $stmt = 'SELECT COUNT(*) FROM ' . $table;
my @params = ();
if (defined $delta) {
$stmt .= ' WHERE ' .
$db->columnidentifier('delta') . ' = ?';
push(@params,$delta);
}
return $db->db_get_value($stmt,@params);
}
sub update_delta {
my ($subscribernumber,$delta) = @_;
check_table();
my $db = &$get_db();
my $table = $db->tableidentifier($tablename);
my $stmt = 'UPDATE ' . $table . ' SET delta = ?';
my @params = ();
push(@params,$delta);
if (defined $subscribernumber) {
$stmt .= ' WHERE ' .
$db->columnidentifier('country_code') . ' = ?' .
' AND ' . $db->columnidentifier('area_code') . ' = ?' .
' AND ' . $db->columnidentifier('dial_number') . ' = ?';
push(@params,split_subscribernumber($subscribernumber));
}
return $db->db_do($stmt,@params);
}
sub buildrecords_fromrows {
my ($rows,$load_recursive) = @_;
@ -142,7 +219,7 @@ sub buildrecords_fromrows {
# transformations go here ...
if ($load_recursive) {
$record->{_features} = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::findby_subscribernumber(
$record->{_features} = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::findby_subscribernumber(
$record->subscribernumber(),
$load_recursive
);
@ -164,6 +241,31 @@ sub getinsertstatement {
}
sub getupsertstatement {
my ($exists_delta,$new_delta) = @_;
check_table();
my $db = &$get_db();
my $table = $db->tableidentifier($tablename);
my $upsert_stmt = 'INSERT OR REPLACE INTO ' . $table . ' (' .
join(', ',map { local $_ = $_; $_ = $db->columnidentifier($_); $_; } @$expected_fieldnames) . ')';
my @values = ();
foreach my $fieldname (@$expected_fieldnames) {
if ('delta' eq $fieldname) {
my $stmt = 'SELECT \'' . $exists_delta . '\' FROM ' . $table . ' WHERE ' .
$db->columnidentifier('country_code') . ' = ?' .
' AND ' . $db->columnidentifier('area_code') . ' = ?' .
' AND ' . $db->columnidentifier('dial_number') . ' = ?';
push(@values,'COALESCE((' . $stmt . '), \'' . $new_delta . '\')');
} else {
push(@values,'?');
}
}
$upsert_stmt .= ' VALUES (' . join(',',@values) . ')';
return $upsert_stmt;
}
sub gettablename {
return $tablename;

@ -31,7 +31,7 @@ OptionValues ::= OptionValue+ action => _build_setoptio
SubscriberNumber ~ [0-9]+
OptionName ~ [-a-zA-Z_0-9]+
OptionValue ~ [-a-zA-Z_0-9 \t]+
OptionValue ~ [-a-zA-Z_0-9 \t#*]+
whitespace ~ [\s]+
:discard ~ whitespace
__GRAMMAR__

@ -40,15 +40,16 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::FileProcessors::LnpDefi
use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw(
get_import_db
destroy_dbs
destroy_all_dbs
);
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp qw();
use NGCP::BulkProcessor::Array qw(removeduplicates);
use NGCP::BulkProcessor::Utils qw(threadid);
require Exporter;
our @ISA = qw(Exporter);
@ -62,8 +63,8 @@ sub import_features_define {
my ($file) = @_;
# create tables:
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::create_table(1);
$result &= NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem::create_table(1);
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::create_table(1);
$result &= NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::create_table(1);
# checks, e.g. other table must be present:
# ..none..
@ -73,8 +74,10 @@ sub import_features_define {
$importer->stoponparseerrors(!$dry);
# launch:
destroy_dbs(); #close all db connections before forking..
return $result && $importer->process($file,sub {
destroy_all_dbs(); #close all db connections before forking..
return $result && $importer->process(
file => $file,
process_code => sub {
my ($context,$rows,$row_offset) = @_;
my $rownum = $row_offset;
my @featureoption_rows = ();
@ -126,7 +129,8 @@ sub import_features_define {
}
}
return 1;
}, sub {
},
init_process_context_code => sub {
my ($context)= @_;
if (not $importer->parselines()) {
eval {
@ -137,19 +141,22 @@ sub import_features_define {
}
}
$context->{db} = &get_import_db(); # keep ref count low..
}, sub {
},
uninit_process_context_code => sub {
my ($context)= @_;
undef $context->{db};
destroy_dbs();
},$import_multithreading);
destroy_all_dbs();
},
multithreading => $import_multithreading
);
}
sub _insert_featureoption_rows {
my ($context,$featureoption_rows) = @_;
$context->{db}->db_do_begin(
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::getinsertstatement($ignore_options_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::gettablename(),
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::getinsertstatement($ignore_options_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::gettablename(),
#lock - $import_multithreading
);
$context->{db}->db_do_rowblock($featureoption_rows);
@ -159,8 +166,8 @@ sub _insert_featureoption_rows {
sub _insert_featureoptionsetitem_rows {
my ($context,$featureoptionsetitem_rows) = @_;
$context->{db}->db_do_begin(
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem::getinsertstatement($ignore_setoptionitems_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem::gettablename(),
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::getinsertstatement($ignore_setoptionitems_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::gettablename(),
#lock
);
$context->{db}->db_do_rowblock($featureoptionsetitem_rows);
@ -171,57 +178,74 @@ sub import_subscriber_define {
my ($file) = @_;
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber::create_table(1);
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::create_table(0);
$result &= _import_subscriber_define_checks($file);
my $importer = NGCP::BulkProcessor::Projects::Migration::IPGallery::FileProcessors::SubscriberDefineFile->new($subscriber_define_import_numofthreads);
destroy_dbs(); #close all db connections before forking..
return $result && $importer->process($file,sub {
my $upsert = _import_subscriber_define_reset_delta();
destroy_all_dbs(); #close all db connections before forking..
return $result && $importer->process(
file => $file,
process_code => sub {
my ($context,$rows,$row_offset) = @_;
my $rownum = $row_offset;
my @subscriber_rows = ();
foreach my $row (@$rows) {
$rownum++;
my $record = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber->new($row);
my $record = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber->new($row);
next if 'None' eq $record->{rgw_fqdn};
if ($record->{dial_number} =~ $subscribernumer_exclude_pattern) {
if ($record->{dial_number} =~ $subscribernumer_exclude_exception_pattern) {
processing_info($context->{tid},'record ' . $rownum . ' - exclude exception pattern match: ' . $record->{dial_number},getlogger(__PACKAGE__));
push(@subscriber_rows,$row) if _import_subscriber_define_referential_checks($context,$record,$rownum);
next unless _import_subscriber_define_referential_checks($context,$record,$rownum);
} else {
processing_info($context->{tid},'record ' . $rownum . ' - skipped, exclude pattern match: ' . $record->{dial_number},getlogger(__PACKAGE__));
next;
}
} else {
push(@subscriber_rows,$row) if _import_subscriber_define_referential_checks($context,$record,$rownum);
next unless _import_subscriber_define_referential_checks($context,$record,$rownum);
}
my @subscriber_row = @$row;
if ($context->{upsert}) {
push(@subscriber_row,$record->{country_code},$record->{area_code},$record->{dial_number});
} else {
push(@subscriber_row,$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::added_delta);
}
push(@subscriber_rows,\@subscriber_row);
}
if ((scalar @$rows) > 0) {
if ((scalar @subscriber_rows) > 0) {
if ($dry) {
eval { _insert_subscriber_rows($context,$rows); };
eval { _insert_subscriber_rows($context,\@subscriber_rows); };
} else {
_insert_subscriber_rows($context,$rows);
_insert_subscriber_rows($context,\@subscriber_rows);
}
}
return 1;
}, sub {
},
init_process_context_code => sub {
my ($context)= @_;
$context->{db} = &get_import_db(); # keep ref count low..
}, sub {
$context->{upsert} = $upsert;
},
uninit_process_context_code => sub {
my ($context)= @_;
undef $context->{db};
destroy_dbs();
}, $import_multithreading);
destroy_all_dbs();
},
multithreading => $import_multithreading
);
}
sub _import_subscriber_define_referential_checks {
my ($context,$record,$rownum) = @_;
my $result = 0;
if (NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::countby_subscribernumber($record->subscribernumber()) > 0) {
if (NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::countby_subscribernumber($record->subscribernumber()) > 0) {
$result = 1;
} else {
if ($dry) {
@ -239,7 +263,7 @@ sub _import_subscriber_define_checks {
my $result = 1;
my $optioncount = 0;
eval {
$optioncount = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::countby_subscribernumber();
$optioncount = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::countby_subscribernumber();
};
if ($@ or $optioncount == 0) {
fileprocessingerror($file,'please import subscriber features first',getlogger(__PACKAGE__));
@ -248,11 +272,27 @@ sub _import_subscriber_define_checks {
return $result;
}
sub _import_subscriber_define_reset_delta {
if (NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_subscribernumber() > 0) {
processing_info(threadid(),'resetting subscriber delta of ' .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::update_delta(undef,
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::deleted_delta) .
' records',getlogger(__PACKAGE__));
return 1;
}
return 0;
}
sub _insert_subscriber_rows {
my ($context,$subscriber_rows) = @_;
$context->{db}->db_do_begin(
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber::getinsertstatement($ignore_subscriber_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber::gettablename(),
($context->{upsert} ?
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::getupsertstatement(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::updated_delta,
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::added_delta
)
: NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::getinsertstatement($ignore_subscriber_unique)),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::gettablename(),
#lock
);
$context->{db}->db_do_rowblock($subscriber_rows);
@ -263,12 +303,14 @@ sub import_lnp_define {
my ($file) = @_;
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp::create_table(1);
my $result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::create_table(1);
my $importer = NGCP::BulkProcessor::Projects::Migration::IPGallery::FileProcessors::LnpDefineFile->new($lnp_define_import_numofthreads);
destroy_dbs(); #close all db connections before forking..
return $result && $importer->process($file,sub {
destroy_all_dbs(); #close all db connections before forking..
return $result && $importer->process(
file => $file,
process_code => sub {
my ($context,$rows,$row_offset) = @_;
my $rownum = $row_offset;
my @lnp_rows = ();
@ -289,22 +331,26 @@ sub import_lnp_define {
}
return 1;
}, sub {
},
init_process_context_code => sub {
my ($context)= @_;
$context->{db} = &get_import_db(); # keep ref count low..
}, sub {
},
uninit_process_context_code => sub {
my ($context)= @_;
undef $context->{db};
destroy_dbs();
}, $import_multithreading);
destroy_all_dbs();
},
multithreading => $import_multithreading
);
}
sub _insert_lnp_rows {
my ($context,$lnp_rows) = @_;
$context->{db}->db_do_begin(
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp::getinsertstatement($ignore_lnp_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp::gettablename(),
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::getinsertstatement($ignore_lnp_unique),
#NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::gettablename(),
#lock
);
$context->{db}->db_do_rowblock($lnp_rows);

@ -26,7 +26,7 @@ use NGCP::BulkProcessor::SqlConnectors::SQLiteDB qw(
#use NGCP::BulkProcessor::SqlConnectors::SQLServerDB;
#use NGCP::BulkProcessor::RestConnectors::NGCPRestApi;
use NGCP::BulkProcessor::SqlRecord qw(cleartableinfo);
use NGCP::BulkProcessor::SqlProcessor qw(cleartableinfo);
require Exporter;
our @ISA = qw(Exporter);
@ -35,6 +35,7 @@ our @EXPORT_OK = qw(
import_db_tableidentifier
destroy_dbs
destroy_all_dbs
);
# thread connector pools:
@ -80,4 +81,9 @@ sub destroy_dbs {
}
sub destroy_all_dbs() {
destroy_dbs();
NGCP::BulkProcessor::ConnectorPool::destroy_dbs();
}
1;

@ -67,6 +67,7 @@ our @EXPORT_OK = qw(
$lnp_define_import_numofthreads
$ignore_lnp_unique
$stats_record_list_limit
);
our $defaultconfig = 'config.cfg';
@ -98,6 +99,9 @@ our $lnp_define_filename = undef;
our $lnp_define_import_numofthreads = $cpucount;
our $ignore_lnp_unique = 1;
our $stats_record_list_limit = 1000;
sub update_settings {
my ($data,$configfile) = @_;

@ -23,6 +23,7 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::Settings qw(
$features_define_filename
$subscriber_define_filename
$lnp_define_filename
$stats_record_list_limit
);
use NGCP::BulkProcessor::Logging qw(
init_log
@ -58,8 +59,7 @@ use NGCP::BulkProcessor::Mail qw(
use NGCP::BulkProcessor::SqlConnectors::CSVDB qw(cleanupcvsdirs);
use NGCP::BulkProcessor::SqlConnectors::SQLiteDB qw(cleanupdbfiles);
use NGCP::BulkProcessor::ConnectorPool qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool qw(destroy_all_dbs);
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Import qw(
import_features_define
@ -67,10 +67,10 @@ use NGCP::BulkProcessor::Projects::Migration::IPGallery::Import qw(
import_lnp_define
);
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber qw();
use NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp qw();
scripterror(getscriptpath() . ' already running',getlogger(getscriptpath())) unless flock DATA, LOCK_EX | LOCK_NB; # not tested on windows yet
@ -85,6 +85,8 @@ my $import_features_define_task_opt = 'import_features';
push(@TASK_OPTS,$import_features_define_task_opt);
my $import_subscriber_define_task_opt = 'import_subscriber';
push(@TASK_OPTS,$import_subscriber_define_task_opt);
my $import_truncate_subscriber_task_opt = 'truncate_subscriber';
push(@TASK_OPTS,$import_truncate_subscriber_task_opt);
my $import_lnp_define_task_opt = 'import_lnp';
push(@TASK_OPTS,$import_lnp_define_task_opt);
@ -137,6 +139,8 @@ sub main() {
$result = import_features_define_task(\@messages) if taskinfo($import_features_define_task_opt,$result);
} elsif (lc($import_subscriber_define_task_opt) eq lc($task)) {
$result = import_subscriber_define_task(\@messages) if taskinfo($import_subscriber_define_task_opt,$result);
} elsif (lc($import_truncate_subscriber_task_opt) eq lc($task)) {
$result = import_truncate_subscriber_task(\@messages) if taskinfo($import_truncate_subscriber_task_opt,$result);
} elsif (lc($import_lnp_define_task_opt) eq lc($task)) {
$result = import_lnp_define_task(\@messages) if taskinfo($import_lnp_define_task_opt,$result);
} elsif (lc('blah') eq lc($task)) {
@ -203,16 +207,19 @@ sub import_features_define_task {
$result = import_features_define($features_define_filename);
};
my $err = $@;
my $stats = ' feature option: ' .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOption::countby_subscribernumber() . ' rows';
$stats .= "\n feature set option items: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::FeatureOptionSetItem::countby_subscribernumber_option() . ' rows';
my $stats = '';
eval {
$stats .= "\n total feature option records: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOption::countby_subscribernumber() . ' rows';
$stats .= "\n total feature set option item records: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::FeatureOptionSetItem::countby_subscribernumber_option() . ' rows';
};
if ($err or !$result) {
push(@$messages,"importing subscriber features incomplete\n$stats");
push(@$messages,"importing subscriber features incomplete$stats");
} else {
push(@$messages,"importing subscriber features completed\n$stats");
push(@$messages,"importing subscriber features completed$stats");
}
destroy_dbs(); #every task should leave with closed connections.
destroy_all_dbs(); #every task should leave with closed connections.
return $result;
}
@ -225,50 +232,106 @@ sub import_subscriber_define_task {
$result = import_subscriber_define($subscriber_define_filename);
};
my $err = $@;
my $stats = ' subscriber: ' .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Subscriber::countby_subscribernumber() . ' rows';
my $stats = '';
eval {
$stats .= "\n total subscriber records: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_subscribernumber() . ' rows';
my $added_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::added_delta
);
$stats .= "\n new: $added_count rows";
if ($added_count <= $stats_record_list_limit) {
foreach my $record (@{NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::findby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::added_delta
)}) {
$stats .= "\n " . $record->{dial_number};
}
}
my $existing_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::updated_delta
);
$stats .= "\n existing: $existing_count rows";
if ($existing_count <= $stats_record_list_limit) {
foreach my $record (@{NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::findby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::updated_delta
)}) {
$stats .= "\n " . $record->{dial_number};
}
}
my $deleted_count = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::deleted_delta
);
$stats .= "\n removed: $deleted_count rows";
if ($deleted_count <= $stats_record_list_limit) {
foreach my $record (@{NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::findby_delta(
$NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::deleted_delta
)}) {
$stats .= "\n " . $record->{dial_number};
}
}
};
if ($err or !$result) {
push(@$messages,"importing subscribers incomplete\n$stats");
push(@$messages,"importing subscribers incomplete$stats");
} else {
push(@$messages,"importing subscribers completed\n$stats");
push(@$messages,"importing subscribers completed$stats");
}
destroy_dbs(); #every task should leave with closed connections.
destroy_all_dbs(); #every task should leave with closed connections.
return $result;
}
sub import_lnp_define_task {
sub import_truncate_subscriber_task {
my ($messages) = @_;
my $result = 0;
eval {
$result = import_lnp_define($lnp_define_filename);
$result = NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::create_table(1);
};
my $err = $@;
my $stats = ' lnp numbers: ' .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp::countby_lrncode_portednumber() . ' rows';
$stats .= "\n lrn codes: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::Lnp::count_lrncodes();
my $stats = '';
eval {
$stats .= "\n total subscriber records: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Subscriber::countby_subscribernumber() . ' rows';
};
if ($err or !$result) {
push(@$messages,"importing lnp numbers incomplete\n$stats");
push(@$messages,"truncating imported subscribers incomplete$stats");
} else {
push(@$messages,"importing lnp numbers\n$stats");
push(@$messages,"truncating imported subscribers completed$stats");
}
destroy_dbs(); #every task should leave with closed connections.
destroy_all_dbs(); #every task should leave with closed connections.
return $result;
}
sub destroy_dbs() {
NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool::destroy_dbs();
NGCP::BulkProcessor::ConnectorPool::destroy_dbs();
sub import_lnp_define_task {
my ($messages) = @_;
my $result = 0;
eval {
$result = import_lnp_define($lnp_define_filename);
};
my $err = $@;
my $stats = '';
eval {
$stats .= "\n total lnp number records: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::countby_lrncode_portednumber() . ' rows';
$stats .= "\n total lrn codes: " .
NGCP::BulkProcessor::Projects::Migration::IPGallery::Dao::import::Lnp::count_lrncodes();
};
if ($err or !$result) {
push(@$messages,"importing lnp numbers incomplete$stats");
} else {
push(@$messages,"importing lnp numbers$stats");
}
destroy_all_dbs(); #every task should leave with closed connections.
return $result;
}
#END {
# # this should not be required explicitly, but prevents Log4Perl's
# # "rootlogger not initialized error upon exit..
# NGCP::BulkProcessor::Projects::Migration::IPGallery::ProjectConnectorPool::destroy_dbs();
# NGCP::BulkProcessor::ConnectorPool::destroy_dbs();
# destroy_all_dbs
#}
__DATA__

File diff suppressed because it is too large Load Diff

File diff suppressed because it is too large Load Diff
Loading…
Cancel
Save