|
@@ -21,6 +21,7 @@ import posixpath
|
|
|
|
|
|
|
|
from django.core import management
|
|
from django.core import management
|
|
|
from django.core.management.base import NoArgsCommand
|
|
from django.core.management.base import NoArgsCommand
|
|
|
|
|
+from django.utils.translation import ugettext as _
|
|
|
|
|
|
|
|
from hadoop import cluster
|
|
from hadoop import cluster
|
|
|
|
|
|
|
@@ -36,28 +37,27 @@ class Command(NoArgsCommand):
|
|
|
remote_fs = cluster.get_hdfs()
|
|
remote_fs = cluster.get_hdfs()
|
|
|
remote_dir = Workflow.objects.create_data_dir(remote_fs)
|
|
remote_dir = Workflow.objects.create_data_dir(remote_fs)
|
|
|
|
|
|
|
|
- # Copy sample binaries
|
|
|
|
|
|
|
+ # Copy examples binaries
|
|
|
for demo in ('lib', 'pig'):
|
|
for demo in ('lib', 'pig'):
|
|
|
local_dir = posixpath.join(LOCAL_SAMPLE_DIR.get(), demo)
|
|
local_dir = posixpath.join(LOCAL_SAMPLE_DIR.get(), demo)
|
|
|
remote_data_dir = posixpath.join(remote_dir, demo)
|
|
remote_data_dir = posixpath.join(remote_dir, demo)
|
|
|
- LOG.info('Copying workflows %s to %s\n' % (local_dir, remote_data_dir))
|
|
|
|
|
|
|
+ LOG.info(_('Copying examples %(local_dir)s to %(remote_data_dir)s\n') % {
|
|
|
|
|
+ 'local_dir': local_dir, 'remote_data_dir': remote_data_dir})
|
|
|
copy_dir(local_dir, remote_fs, remote_data_dir)
|
|
copy_dir(local_dir, remote_fs, remote_data_dir)
|
|
|
|
|
|
|
|
# Copy sample data
|
|
# Copy sample data
|
|
|
local_dir = LOCAL_SAMPLE_DATA_DIR.get()
|
|
local_dir = LOCAL_SAMPLE_DATA_DIR.get()
|
|
|
remote_data_dir = posixpath.join(remote_dir, 'data')
|
|
remote_data_dir = posixpath.join(remote_dir, 'data')
|
|
|
- LOG.info('Copying data %s to %s\n' % (local_dir, remote_data_dir))
|
|
|
|
|
|
|
+ LOG.info(_('Copying data %(local_dir)s to %(remote_data_dir)s\n') % {
|
|
|
|
|
+ 'local_dir': local_dir, 'remote_data_dir': remote_data_dir})
|
|
|
copy_dir(local_dir, remote_fs, remote_data_dir)
|
|
copy_dir(local_dir, remote_fs, remote_data_dir)
|
|
|
|
|
|
|
|
# Load jobs
|
|
# Load jobs
|
|
|
management.call_command('loaddata', 'apps/oozie/src/oozie/fixtures/initial_data.json', verbosity=2)
|
|
management.call_command('loaddata', 'apps/oozie/src/oozie/fixtures/initial_data.json', verbosity=2)
|
|
|
|
|
|
|
|
- def has_been_setup(self):
|
|
|
|
|
- return False
|
|
|
|
|
-
|
|
|
|
|
|
|
|
|
|
def copy_dir(local_dir, remote_fs, remote_dir, mode=755):
|
|
def copy_dir(local_dir, remote_fs, remote_dir, mode=755):
|
|
|
- remote_fs.mkdir(remote_dir, mode=mode)
|
|
|
|
|
|
|
+ remote_fs.do_as_user(remote_fs.DEFAULT_USER, remote_fs.mkdir, remote_dir, mode=mode)
|
|
|
|
|
|
|
|
for f in os.listdir(local_dir):
|
|
for f in os.listdir(local_dir):
|
|
|
local_src = os.path.join(local_dir, f)
|
|
local_src = os.path.join(local_dir, f)
|
|
@@ -69,21 +69,25 @@ CHUNK_SIZE = 1024 * 1024
|
|
|
|
|
|
|
|
def copy_file(local_src, remote_fs, remote_dst):
|
|
def copy_file(local_src, remote_fs, remote_dst):
|
|
|
if remote_fs.exists(remote_dst):
|
|
if remote_fs.exists(remote_dst):
|
|
|
- LOG.info('%s already exists. Skipping.' % remote_dst)
|
|
|
|
|
|
|
+ LOG.info(_('%(remote_dst)s already exists. Skipping.') % {'remote_dst': remote_dst})
|
|
|
return
|
|
return
|
|
|
else:
|
|
else:
|
|
|
- LOG.info('%s does not exist. trying to copy' % remote_dst)
|
|
|
|
|
|
|
+ LOG.info(_('%(remote_dst)s does not exist. Trying to copy') % {'remote_dst': remote_dst})
|
|
|
|
|
|
|
|
if os.path.isfile(local_src):
|
|
if os.path.isfile(local_src):
|
|
|
src = file(local_src)
|
|
src = file(local_src)
|
|
|
try:
|
|
try:
|
|
|
- remote_fs.create(remote_dst, permission=01755)
|
|
|
|
|
- chunk = src.read(CHUNK_SIZE)
|
|
|
|
|
- while chunk:
|
|
|
|
|
- remote_fs.append(remote_dst, chunk)
|
|
|
|
|
|
|
+ try:
|
|
|
|
|
+ remote_fs.do_as_user(remote_fs.DEFAULT_USER, remote_fs.create, remote_dst, permission=01755)
|
|
|
chunk = src.read(CHUNK_SIZE)
|
|
chunk = src.read(CHUNK_SIZE)
|
|
|
- LOG.info('Copied %s -> %s' % (local_src, remote_dst))
|
|
|
|
|
|
|
+ while chunk:
|
|
|
|
|
+ remote_fs.do_as_user(remote_fs.DEFAULT_USER, remote_fs.append, remote_dst, chunk)
|
|
|
|
|
+ chunk = src.read(CHUNK_SIZE)
|
|
|
|
|
+ LOG.info(_('Copied %s -> %s') % (local_src, remote_dst))
|
|
|
|
|
+ except:
|
|
|
|
|
+ LOG.error(_('Copying %s -> %s failed') % (local_src, remote_dst))
|
|
|
|
|
+ raise
|
|
|
finally:
|
|
finally:
|
|
|
src.close()
|
|
src.close()
|
|
|
else:
|
|
else:
|
|
|
- LOG.info('Skipping %s (not a file)' % local_src)
|
|
|
|
|
|
|
+ LOG.info(_('Skipping %s (not a file)') % local_src)
|