ArielGlenn has uploaded a new change for review. (
https://gerrit.wikimedia.org/r/342846 )
Change subject: scripts to generate a series of checkpoint files for a dump run
manually
......................................................................
scripts to generate a series of checkpoint files for a dump run manually
* script for locking the a wiki dump run and continuously updating lock;
this is a standalone script that can be used from e.g. bash wrappers
* script for generating start/end page ids suitable for batching up
a pile of jobs to produce page-meta-history output files in a reasonable
timeframe, with all jobs finishing up around the same time
* script for generating page meta history content for missing page
ranges for a set of subjobs
These are meant to be used when recovering from a broken run manually
(when you're in a hurry for the run to complete).
Bug: T160507
Change-Id: I10796ae540f04a8b16c604ffc0d7ff681de9d06e
---
A xmldumps-backup/dump_lock.py
A xmldumps-backup/pagerange.py
A xmldumps-backup/runme.sh
3 files changed, 848 insertions(+), 0 deletions(-)
git pull ssh://gerrit.wikimedia.org:29418/operations/dumps
refs/changes/46/342846/1
diff --git a/xmldumps-backup/dump_lock.py b/xmldumps-backup/dump_lock.py
new file mode 100644
index 0000000..c9cd283
--- /dev/null
+++ b/xmldumps-backup/dump_lock.py
@@ -0,0 +1,144 @@
+"""
+get a lock on the specified wiki and date,
+touch it every so often so it doesn't get stale,
+listen for a hup and if there is one, remove
+the lock and go home.
+this is designed for use be shell scripts
+that run standalone sets of jobs and get the
+lock at the beginning, releaseing it for the end.
+"""
+
+
+import sys
+import signal
+import getopt
+import time
+from dumps.WikiDump import Locker
+from dumps.WikiDump import Config
+from dumps.WikiDump import Wiki
+
+
+class StandaloneLocker(object):
+ """
+ lock a wiki run for a given date
+ return True if lock acquired, False otherwise
+ on HUP, removes lock and exits with 0
+ """
+ def __init__(self, wikiname, configfile, date):
+ self.wiki = Wiki(Config(configfile), wikiname)
+ self.wiki.set_date(date)
+ self.locker = None
+
+ def get_lock(self):
+ """
+ get the lock on the specified wiki
+ or return False;
+ this also automatically starts up
+ a watchdog that touches the lock file
+ every so often.
+ """
+ try:
+ locker = Locker(self.wiki, self.wiki.date)
+ locker.lock()
+ self.locker = locker
+ return True
+ except Exception as ex:
+ return False
+
+ def handle_hup(self, signo, dummy_frame):
+ """
+ ignore any more hups
+ stop the lock refresher and remove the lock
+ go bye bye
+ """
+ signal.signal(signal.SIGHUP, signal.SIG_IGN)
+ lockfiles = self.locker.is_locked()
+ self.locker.unlock(lockfiles, owner=True)
+ sys.exit(0)
+
+
+def usage(message=None):
+ '''
+ display a helpful usage message with
+ an optional introductory message first
+ '''
+ if message is not None:
+ sys.stderr.write(message)
+ sys.stderr.write("\n")
+ usage_message = """
+Usage: dump_lock.py --wiki <wikiname> --date <YYYYMMDD>
+ --configfile <path> [--verbose] [--help]
+
+--wiki (-w): name of db of wiki for which to get lock
+--date (-j): date of run to lock
+--configfile (-c): path to config file
+--verbose (-v): display messages about what the script is doing
+--help (-h): display this help message
+"""
+ sys.stderr.write(usage_message)
+ sys.exit(1)
+
+
+def get_args():
+ """
+ get and validate args, and return them
+ """
+ wiki = None
+ date = None
+ configfile = None
+ verbose = False
+ try:
+ (options, remainder) = getopt.gnu_getopt(
+ sys.argv[1:], "c:w:d:vh", ["configfile=", "wiki=", "date=",
+ "verbose", "help"])
+ except getopt.GetoptError as err:
+ usage("Unknown option specified: " + str(err))
+
+ for (opt, val) in options:
+ if opt in ["-c", "--configfile"]:
+ configfile = val
+ elif opt in ["-w", "--wiki"]:
+ wiki = val
+ elif opt in ["-j", "--date"]:
+ date = val
+ elif opt in ["-v", "--verbose"]:
+ verbose = True
+ elif opt in ["-h", "--help"]:
+ usage("Help for this script")
+
+ if not date.isdigit or len(date) != 8:
+ usage("'date' must be in format YYYYMMDD")
+
+ if len(remainder) > 0:
+ usage("Unknown option specified")
+
+ if not wiki or not date or not configfile:
+ usage("Arguments 'wiki', 'date' and 'configfile' must be set")
+ return(wiki, date, configfile, verbose)
+
+
+def do_main():
+ """entry point:
+ get args, try to lock, nonzero exit code on failure
+ """
+ wiki, date, configfile, verbose = get_args()
+
+ stlocker = StandaloneLocker(wiki, configfile, date)
+ signal.signal(signal.SIGHUP, stlocker.handle_hup)
+ if not stlocker.get_lock():
+ if verbose:
+ sys.stderr.write("Failed to get lock for {wiki}, {date}\n".format(
+ wiki=wiki, date=date))
+ sys.exit(1)
+
+ if verbose:
+ sys.stderr.write("Lock acquired for {wiki}, {date}\n".format(
+ wiki=wiki, date=date))
+
+ # sit here mostly lazing around until HUP time
+ while True:
+ time.sleep(60)
+
+
+if __name__ == '__main__':
+ do_main()
diff --git a/xmldumps-backup/pagerange.py b/xmldumps-backup/pagerange.py
new file mode 100644
index 0000000..35ec8c1
--- /dev/null
+++ b/xmldumps-backup/pagerange.py
@@ -0,0 +1,493 @@
+"""
+given either number of jobs desired, or number
+of revs (approx) per job desired, generate a list
+of sequential page ranges that could be dumped to
+produce the required number of jobs all taking
+about the same length of time, or to produce
+the requested number of revisions in each output,
+likewise taking (roughly) the same amount of
+time.
+this lets us break up full content dumps into
+a number of smaller jobs and run them in batches,
+rerunning specific files as needed.
+"""
+
+
+import getopt
+import sys
+import time
+import json
+
+from dumps.WikiDump import Config, Wiki
+from dumps.utils import DbServerInfo
+import xmlstreams
+
+
+def get_count_from_output(sqloutput):
+ """
+ given sql query output,
+ return the single integer value we expect from it
+ """
+ lines = sqloutput.splitlines()
+ if lines and lines[1]:
+ if not lines[1].isdigit():
+ return None # probably NULL or missing table
+ return int(lines[1])
+ return None
+
+
+def get_estimate_from_output(sqloutput):
+ """
+ dig out the estimated number or rows form sql output
+ and return it
+ """
+ lines = sqloutput.splitlines()
+ if lines and lines[1]:
+ fields = lines[1].split()
+ # id | select_type | table | type | possible_keys | key | key_len |
ref | rows
+ if not fields[8].isdigit():
+ print "unexpected output from sql query, giving up:"
+ print sqloutput
+ sys.exit(1)
+ return int(fields[8])
+ return None
+
+
+def adjust(count, adj, sign):
+ """
+ adjust the count by the adjustment amount in
+ the direction of the sign
+ """
+ if sign < 0:
+ count -= adj
+ else:
+ count += adj
+ return count
+
+
+class QueryRunner(object):
+ """
+ runs various db queries related to page, revision count, etc.
+ """
+ def __init__(self, dbname, config, verbose=False):
+ self.dbname = dbname
+ self.config = config
+ self.verbose = verbose
+ self.wiki = Wiki(self.config, self.dbname)
+ self.db_info = DbServerInfo(self.wiki, self.dbname)
+
+ def get_max_id(self, idtype):
+ """
+ get and return the max (rev or page) id
+ """
+ if idtype == 'page':
+ return xmlstreams.get_max_id(self.config, self.dbname, 'page_id',
'page')
+ elif idtype == 'rev':
+ return xmlstreams.get_max_id(self.config, self.dbname, 'rev_id',
'revision')
+ else:
+ return None
+
+ def get_count(self, page_start, page_end):
+ """
+ get the number of revisions for the pages starting
+ from page_start and ending with page_end
+ and return it
+ """
+ query = ("select count(rev_id) from revision where "
+ "rev_page >= {start} and rev_page < {end}".format(
+ start=page_start, end=page_end))
+ queryout = self.run_sql_query(query)
+ if queryout is None:
+ print "unexpected output from sql query, giving up:"
+ print query, queryout
+ sys.exit(1)
+
+ revcount = get_count_from_output(queryout)
+ if revcount is None:
+ print "unexpected output from sql query, giving up:"
+ print query, queryout
+ sys.exit(1)
+ return revcount
+
+ def get_estimate(self, page_start, page_end):
+ """
+ get estimate of number of revisions (via explain)
+ for page range from page_start to page_end
+ and return it
+ """
+ query = ("explain select count(rev_id) from revision where "
+ "rev_page >= {start} and rev_page < {end}".format(
+ start=page_start, end=page_end))
+ queryout = self.run_sql_query(query)
+ if queryout is None:
+ print "unexpected output from sql query, giving up:"
+ print query, queryout
+ sys.exit(1)
+ return get_estimate_from_output(queryout)
+
+ def run_sql_query(self, query, maxretries=3):
+ """
+ run the supplied sql query, retrying up to
+ maxtretries times, with sleep of 5 secs in between
+ this will recheck the db config in case a db server
+ has been removed from the pool suddenly
+ """
+ results = None
+ retries = 0
+ while results is None and retries < maxretries:
+ retries = retries + 1
+ results = self.db_info.run_sql_and_get_output(query)
+ if results is None:
+ time.sleep(5)
+ # get db config again in case something's changed
+ self.db_info = DbServerInfo(self.wiki, self.dbname)
+ continue
+ return results
+ return results
+
+
+class PageRange(object):
+ '''
+ Methods for getting number of revisions for a page range.
+ We use this for splitting up history runs into small chunks to be run in
+ parallel, with each job taking roughly the same length of time.
+ '''
+
+ def __init__(self, runner, verbose=False):
+ '''
+ Arguments:
+ dbname -- the name of the database we are dumping
+ config -- this is the general config context used by the runner class
+ '''
+ self.runner = runner
+ self.verbose = verbose
+
+ self.total_pages = runner.get_max_id('page')
+ self.total_revs = runner.get_max_id('rev')
+
+ def get_page_ranges_for_jobs(self, numjobs):
+ '''
+ get and return list of tuples consisting of page id start and end to
be passed to
+ self.numjobs jobs which should run in approximately the same length of
time
+ for full history dumps, which is all we care about, really
+ numjobs -- number of jobs to be run in parallel, and so number of page
ranges
+ we want to produce for these parallel runs
+ '''
+
+ ranges = []
+ page_start = 1
+ numrevs = self.total_revs / numjobs + 1
+ # actually this is ok for the start but it varies right afterwards
+ # interval = ((self.total_pages - page_start)/numjobs_left) + 1
+ prevguess = 0
+ for jobnum in range(1, numjobs + 1):
+ if jobnum == numjobs:
+ # last job, don't bother searching. just append up to max page
id
+ ranges.append((page_start, self.total_pages))
+ break
+ # this is wrong too, we need to get it passed or something
+# prevguess = min(interval*(jobnum+1), self.total_pages)
+ # fixme here the interval*jobnum can't be right
+# (start, end) = self.get_page_range(page_start, numrevs,
+# interval*jobnum, prevguess)
+ numjobs_left = numjobs - jobnum + 1
+ interval = ((self.total_pages - page_start) / numjobs_left) + 1
+ (start, end) = self.get_page_range(page_start, numrevs,
+ page_start + interval,
prevguess)
+ page_start = end + 1
+ prevguess = end
+ if end > self.total_pages:
+ end = self.total_pages
+ ranges.append((start, end))
+ if page_start > self.total_pages:
+ break
+ return ranges
+
+ def get_page_ranges_for_revs(self, page_start, page_end, numrevs):
+ '''
+ get and return list of tuples consisting of page id start and end
+ which should each, if dumped (full history content dumps) contain about
+ the specified number of revisions, and thus run in something close
+ to the same time
+ numrevs -- number of revisions (approx) for each page range to
contain
+ page_start -- don't start at page 1, start at this page instead
+ page_end -- don't end with last page, end at this page instead
+ '''
+
+ ranges = []
+ if not page_start:
+ page_start = 1
+ if not page_end:
+ page_end = self.total_pages
+ # actually this is ok for the start but it varies right afterwards
+ # interval = ((self.total_pages - page_start)/numjobs_left) + 1
+ prevguess = 0
+ if page_start == 1 and page_end == self.total_pages:
+ numjobs = self.total_revs / numrevs + 1
+ else:
+ estimate = self.runner.get_estimate(page_start, page_end)
+ revs_for_range = self.get_revcount(page_start, page_end, estimate)
+ numjobs = revs_for_range / numrevs + 1
+ for jobnum in range(1, numjobs + 1):
+ if jobnum == numjobs:
+ # last job, don't bother searching. just append up to max page
id
+ ranges.append((page_start, page_end))
+ break
+ # this is wrong too, we need to get it passed or something
+# prevguess = min(interval*(jobnum+1), self.total_pages)
+ # fixme here the interval*jobnum can't be right
+# (start, end) = self.get_page_range(page_start, numrevs,
+# interval*jobnum, prevguess)
+ numjobs_left = numjobs - jobnum + 1
+ interval = (page_end - page_start) / numjobs_left + 1
+ (start, end) = self.get_page_range(page_start, numrevs,
+ page_start + interval,
prevguess)
+ page_start = end + 1
+ prevguess = end
+ if end > page_end:
+ end = page_end
+ ranges.append((start, end))
+ if page_start > page_end:
+ break
+ return ranges
+
+ def get_revcount(self, page_start, page_end, estimate):
+ """
+ for the given page range, get the number of revisions,
+ running multiple queries so we don't make the servers sad
+ the number of queries is based on the estimated number
+ of revs passed in ("estimate"), where we try not to run
+ a query that will result in more than 50k revs being
+ counted. Key word being "try". IN some cases a page
+ all by itself may have hundreds of thousands of revs,
+ whaddya gonna do
+ """
+ total = 0
+ maxtodo = 50000
+
+ runstodo = estimate / maxtodo + 1
+ step = (page_end - page_start) / runstodo
+ ends = range(page_start, page_end, step)
+
+ if ends[-1] != page_end:
+ ends.append(page_end)
+ interval_start = ends.pop(0)
+
+ for interval_end in ends:
+ count = self.runner.get_count(interval_start, interval_end)
+ interval_start = interval_end + 1
+ total += count
+ return total
+
+ def get_page_range(self, page_start, numrevs, badguess, prevguess):
+ """
+ given starting page, number of revisions desired for the page
+ range, the current end of range guess and the previous guess,
+ progressively narrow down the endpoint til we get a value
+ that is "close enough" and return the page range
+ as tuple (start, end)
+ """
+ if self.verbose:
+ print ("get_page_range called with page_start", page_start,
+ "numrevs", numrevs, "badguess", badguess, "prevguess",
prevguess)
+ interval_start = page_start
+ interval_end = badguess
+ revcount = 0
+ while True:
+ if self.verbose:
+ print ("page range loop, start page_id",
+ interval_start, "end page_id:", interval_end)
+
+ estimate = self.runner.get_estimate(interval_start, interval_end)
+ revcount_adj = self.get_revcount(interval_start, interval_end,
estimate)
+ revcount = adjust(revcount, revcount_adj, badguess - prevguess)
+
+ if self.verbose:
+ print ("estimate is", estimate, "revcount is", revcount,
+ "and numrevs is", numrevs)
+
+ interval = abs(prevguess - badguess) / 2
+ if not interval:
+ return (page_start, badguess)
+ prevguess = badguess
+
+ margin = abs(revcount - numrevs)
+ if margin <= 100:
+ return (page_start, badguess)
+ if self.verbose:
+ print "revcount is greater than allowed margin from numrevs"
+
+ badguess = adjust(badguess, interval, numrevs - revcount)
+
+ if badguess < prevguess:
+ interval_start = badguess
+ interval_end = prevguess
+ else:
+ interval_start = prevguess
+ interval_end = badguess
+ if self.verbose:
+ print "new interval:", interval_start, interval_end
+
+
+def usage(message=None):
+ '''
+ display a helpful usage message with
+ an optional introductory message first
+ '''
+ if message is not None:
+ sys.stderr.write(message)
+ sys.stderr.write("\n")
+ usage_message = """
+Usage: pagerange.py --wiki <wikiname> --jobs|--revs <int>
+ [--pagestart <int>] [--pageend <int>]
+ [--pad <int>] [--json]
+ [--configfile <path>] [--verbose] [--help]
+
+--wiki (-w): name of db of wiki for which to run
+--jobs (-j): generate page ranges for this number of jobs
+--revs (-r): generate page ranges for this number of
+ revisions per interval
+--pagestart (-s): page id start for --revs option, default 1
+--pageend (-e): page id end for --revs option, default last
+--configfile (-c): path to config file
+--pad (-p): pad numbers out to specified length with
+ leading zeros for json format
+--json (-J): write results as json formatted output
+--verbose (-v): display messages about what the script is doing
+--help (-h): display this help message
+"""
+ sys.stderr.write(usage_message)
+ sys.exit(1)
+
+
+def check_args(remainder, wiki, revs, jobs):
+ """
+ whine if there are command line arg problems,
+ and exit
+ """
+ if len(remainder) > 0:
+ usage("Unknown option specified")
+
+ if wiki is None:
+ usage("Mandatory option 'wiki' is not specified")
+ if (revs is None and jobs is None) or (revs and jobs):
+ usage("Exactly one of 'revs' or 'jobs' must be specified")
+
+
+def check_int_range_opts(range_opts, names):
+ """
+ make sure they're all nice integers and convert them,
+ or whine and exit
+ """
+
+ for varname in names:
+ if range_opts[varname] is not None:
+ if not range_opts[varname].isdigit:
+ usage("--" + names[varname] + "argument requires a number")
+ else:
+ range_opts[varname] = int(range_opts[varname])
+
+
+def jsonify(pagerange, pad):
+ """
+ given a list of tuples [(startid, endid)...] turn it into
+ a list of hashes [{'start': blah, 'end': blah}...]
+ which can be fed to any json processor
+ optionally (if pad is greater than zero) the
+ ids will be padded out to the specified length
+ with leading zeros
+ """
+ jsonified = []
+ for pair in pagerange:
+ start = str(pair[0])
+ end = str(pair[1])
+ if pad > 0:
+ start = start.zfill(pad)
+ end = end.zfill(pad)
+ jsonified.append({'start': start, 'end': end})
+ return jsonified
+
+
+def do_pageranges(prange, range_opts, pad, jsonfmt):
+ """
+ get and display page ranges for revs or jobs as specified, adding
+ optionally padding to 'pad' places by leading zeros, if 0
+ is passed no padding will be done
+ """
+ if range_opts['jobs']:
+ ranges = prange.get_page_ranges_for_jobs(range_opts['jobs'])
+ # convert ranges into the output we need for the pagesperchunkhistory
config
+ pages_per_job = [page_end - page_start for (page_start, page_end) in
ranges]
+ if jsonfmt:
+ print json.dumps(jsonify(ranges, pad))
+ print json.dumps(pages_per_job, pad)
+ else:
+ print "for {jobs} jobs, have
ranges:".format(jobs=range_opts['jobs'])
+ print ranges
+ print "for {jobs} jobs, have config
setting:".format(jobs=range_opts['jobs'])
+ print pages_per_job
+ else:
+ ranges = prange.get_page_ranges_for_revs(range_opts['start'],
range_opts['end'],
+ range_opts['revs'])
+ if jsonfmt:
+ print json.dumps(jsonify(ranges, pad))
+ else:
+ print ("for start", range_opts['start'], "and end",
+ range_opts['end'] if range_opts['end'] is not None else
"last page")
+ print "for {revs} revs, have
ranges:".format(revs=range_opts['revs'])
+ print ranges
+
+
+def do_main():
+ """
+ main entry point
+ """
+ range_opts = {'jobs': None, 'revs': None, 'start': None, 'end': None}
+ wiki = None
+ configpath = "wikidump.conf"
+ jsonfmt = False
+ verbose = False
+ pad = 0
+ try:
+ (options, remainder) = getopt.gnu_getopt(sys.argv[1:],
"c:w:j:s:e:r:p:vh",
+ ["configfile=", "wiki=",
"jobs=",
+ "pagestart=", "pageend=",
"revs=",
+ "pad=", "json", "verbose",
"help"])
+ except getopt.GetoptError as err:
+ usage("Unknown option specified: " + str(err))
+
+ for (opt, val) in options:
+ if opt in ["-c", "--configfile"]:
+ configpath = val
+ elif opt in ["-w", "--wiki"]:
+ wiki = val
+ elif opt in ["-j", "--jobs"]:
+ range_opts['jobs'] = val
+ elif opt in ["-r", "--revs"]:
+ range_opts['revs'] = val
+ elif opt in ["-e", "--pageend"]:
+ range_opts['end'] = val
+ elif opt in ["-s", "--pagestart"]:
+ range_opts['start'] = val
+ elif opt in ["-p", "--pad"]:
+ try:
+ pad = int(val)
+ except ValueError:
+ usage("'pad' argument must be a number")
+ elif opt in ["-J", "--json"]:
+ jsonfmt = True
+ elif opt in ["-v", "--verbose"]:
+ verbose = True
+ elif opt in ["-h", "--help"]:
+ usage("Help for this script")
+
+ check_int_range_opts(range_opts, {'jobs': 'jobs', 'revs': 'revs',
+ 'end': 'pageend', 'start': 'pagestart'})
+ check_args(remainder, wiki, range_opts['revs'], range_opts['jobs'])
+
+ prange = PageRange(QueryRunner(wiki, Config(configpath), verbose), verbose)
+ do_pageranges(prange, range_opts, pad, jsonfmt)
+
+
+if __name__ == "__main__":
+ do_main()
diff --git a/xmldumps-backup/runme.sh b/xmldumps-backup/runme.sh
new file mode 100644
index 0000000..d9bf20a
--- /dev/null
+++ b/xmldumps-backup/runme.sh
@@ -0,0 +1,211 @@
+#!/bin/bash
+# no error checking, we don't care. if file fails we'll
+# rerun it by hand later
+
+# locks wiki for date, generates a bunch of page ranges based on
+# user input, runs a bunch of jobs in batches to create the
+# specified page-meta-history bz2 output files, unlocks wiki.
+# does NOT: update md5s, status, dumprininfo, symlinks, etc.
+# does NOT: clean up old dumps, remove old files from run
+
+usage() {
+ echo "Usage: $0 --config <pathtofile> --wiki <dbname>"
+ echo " --date <YYYYMMDD> --jobinfo num:num:num,..."
+ echo "[--dryrun] [--verbose]"
+ echo
+ echo " --config path to configuration file for dump generation"
+ echo " --wiki dbname of wiki"
+ echo " --jobinfo partnum:start:end,partnum2:start:end,..."
+ echo " leave start and/or end empty if you want to"
+ echo " start from the beginning or go to last page"
+ echo " --date date of run"
+ echo " --dryrun don't run commands, show what would have been done"
+ echo " --verbose print commands as they are run, etc"
+ exit 1
+}
+
+set_defaults() {
+ CONFIGFILE=""
+ WIKI=""
+ JOBINFO=""
+ DATE=""
+ DRYRUN=""
+ VERBOSE=""
+}
+
+process_opts () {
+ while [ $# -gt 0 ]; do
+ if [ $1 == "--config" ]; then
+ CONFIGFILE="$2"
+ shift; shift;
+ elif [ $1 == "--wiki" ]; then
+ WIKI="$2"
+ shift; shift
+ elif [ $1 == "--jobinfo" ]; then
+ JOBINFO="$2"
+ shift; shift
+ elif [ $1 == "--date" ]; then
+ DATE="$2"
+ shift; shift
+ elif [ $1 == "--dryrun" ]; then
+ DRYRUN="true"
+ shift
+ elif [ $1 == "--verbose" ]; then
+ VERBOSE="true"
+ shift
+ else
+ echo "$0: Unknown option $1"
+ usage
+ fi
+ done
+}
+
+check_opts() {
+ if [ -z "$WIKI" -o -z "$JOBINFO" -o -z "$DATE" -o -z "$CONFIGFILE" ]; then
+ echo "$0: Mandatory options 'wiki', 'jobinfo', 'date' and 'config'
must be specified"
+ usage
+ elif [ ! -f "$CONFIGFILE" ]; then
+ echo "Could not find config file: $CONFIGFILE"
+ echo "Exiting..."
+ exit 1
+ fi
+}
+
+setup_pagerange_args() {
+ # set up the command
+ pagerangeargs=( "$WIKIDUMP_BASE/pagerange.py" )
+ pagerangeargs=( "${pagerangeargs[@]}" "--configfile" "$CONFIGFILE" )
+ if [ -n "$START" ]; then
+ pagerangeargs=( "${pagerangeargs[@]}" "--pagestart" "$START" )
+ fi
+ if [ -n "$END" ]; then
+ pagerangeargs=( "${pagerangeargs[@]}" "--pageend" "$END" )
+ fi
+ pagerangeargs=( "${pagerangeargs[@]}" "--wiki" "$WIKI" )
+ # need json out for our purposes
+ pagerangeargs=( "${pagerangeargs[@]}" "--json" )
+ # 1 million revs per output file is good
+ pagerangeargs=( "${pagerangeargs[@]}" "--revs" "1000000" )
+ # page ids in filenames are always padded out to 9 places
+ pagerangeargs=( "${pagerangeargs[@]}" "--pad" "9" )
+
+ # pipeline args go here
+ grepargs=( "grep" "-v" "DEBUG" )
+ jqargs=( "/usr/bin/jq" "-r" '.[]|.start+":"+.end' )
+}
+
+get_ranges() {
+ if [ -n "$VERBOSE" ]; then
+ echo "/usr/bin/python ${pagerangeargs[@]} | ${grepargs[@]} |
${jqargs[@]}"
+ fi
+ ranges=( $(/usr/bin/python ${pagerangeargs[@]} | ${grepargs[@]} |
${jqargs[@]}) )
+ result=$?
+ if [ $result -ne 0 ]; then
+ echo "Failed to get page ranges, dumping them here"
+ echo "${ranges[@]}"
+ exit 1
+ fi
+}
+
+setup_worker_args() {
+ FILE="$1"
+ # set up the command
+ workerargs=( "$WIKIDUMP_BASE/worker.py" )
+ workerargs=( "${workerargs[@]}" "--configfile" "$CONFIGFILE" )
+ # workerargs=( "${workerargs[@]}" "--log" )
+ workerargs=( "${workerargs[@]}" "--job" "metahistorybz2dump" )
+ # sanity check of date
+ result=`date -d "$DATE"`
+ if [ -z "$result" ]; then
+ echo "bad date given for 'date' arg"
+ exit 1
+ fi
+ workerargs=( "${workerargs[@]}" "--date" "$DATE" )
+ workerargs=( "${workerargs[@]}" "--checkpoint" "$FILE" )
+ workerargs=( "${workerargs[@]}" "$WIKI" )
+}
+
+run_workers() {
+ # run this many workers at once
+ LIMIT=10
+ while :
+ do
+ if [ ${#ranges[*]} -eq 0 ]; then
+ break
+ elif [ ${#ranges[*]} -lt $LIMIT ]; then
+ end=${#ranges[*]}
+ else
+ end=$LIMIT
+ fi
+
+ tenpairs=(${ranges[@]:0:$end})
+
+ wait_pids=()
+ files=()
+ for pagerange in ${tenpairs[@]}; do
+ IFS=: read startpage endpage <<< "$pagerange"
+
outputfile="${WIKI}-${DATE}-pages-meta-history${JOBNUM}.xml-p${startpage}p${endpage}.bz2"
+ setup_worker_args "$outputfile"
+ if [ -n "$DRYRUN" -o -n "$VERBOSE" ]; then
+ echo "${workerargs[@]}"
+ fi
+ if [ -z "$DRYRUN" ]; then
+ /usr/bin/python ${workerargs[@]} &
+ wait_pids+=($!)
+ files+=("$outputfile")
+ fi
+ done
+ i=0
+ for pid in ${wait_pids[*]}; do
+ wait $pid
+ if [ $? -ne 0 ]; then
+ echo "failed to generate" ${files[$i]} "with nonzero exit code"
+ fi
+ ((i++))
+ done
+ ranges=(${ranges[@]:$end})
+ done
+}
+
+lockerup() {
+ if [ -z "$DRYRUN" ]; then
+ /usr/bin/python "$WIKIDUMP_BASE/dump_lock.py" --wiki $WIKI --date
$DATE --configfile $CONFIGFILE &
+ lockerpid=$!
+ sleep 2 # wait a bit, give the process time to finish up if it failed
+ # see if it's still running (which means it got the lock)
+ kill -0 "$lockerpid" >/dev/null 2>&1
+ if [ $? -ne 0 ]; then
+ echo "failed to get lock, exiting"
+ exit 1
+ elif [ -n "$VERBOSE" ]; then
+ echo "got lock"
+ fi
+ fi
+}
+
+cleanup_lock() {
+ if [ -z "$DRYRUN" ]; then
+ if [ -n "$lockerpid" ]; then
+ kill -HUP $lockerpid
+ fi
+ if [ -n "$VERBOSE" ]; then
+ echo "removed lock"
+ fi
+ fi
+}
+
+WIKIDUMP_BASE=`dirname "$0"`
+#DUMPFILESBASE="/mnt/data/xmldatadumps/public"
+DUMPFILESBASE="/home/ariel/dumptesting/dumpruns/public"
+set_defaults
+process_opts "$@"
+check_opts
+IFS=',' read -a array <<< "$JOBINFO"
+lockerup
+for JOB in $JOBINFO; do
+ IFS=: read JOBNUM START END <<< "$JOB"
+ setup_pagerange_args
+ get_ranges
+ run_workers
+done
+cleanup_lock
--
To view, visit https://gerrit.wikimedia.org/r/342846
To unsubscribe, visit https://gerrit.wikimedia.org/r/settings
Gerrit-MessageType: newchange
Gerrit-Change-Id: I10796ae540f04a8b16c604ffc0d7ff681de9d06e
Gerrit-PatchSet: 1
Gerrit-Project: operations/dumps
Gerrit-Branch: master
Gerrit-Owner: ArielGlenn <[email protected]>
_______________________________________________
MediaWiki-commits mailing list
[email protected]
https://lists.wikimedia.org/mailman/listinfo/mediawiki-commits