- Notifications
You must be signed in to change notification settings - Fork 7
Expand file tree
/
Copy pathdatax.py
More file actions
Latest commit
227 lines (189 loc) · 8.79 KB
/
Copy pathdatax.py
File metadata and controls
227 lines (189 loc) · 8.79 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
#!/usr/bin/env python
# -*- coding:utf-8 -*-
importsys
importos
importsignal
importsubprocess
importtime
importre
importsocket
importjson
fromoptparseimportOptionParser
fromoptparseimportOptionGroup
fromstringimportTemplate
importcodecs
importplatform
defisWindows():
returnplatform.system() =='Windows'
DATAX_HOME=os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
DATAX_VERSION='DATAX-OPENSOURCE-3.0'
ifisWindows():
codecs.register(lambdaname: name=='cp65001'andcodecs.lookup('utf-8') orNone)
CLASS_PATH= ("%s/lib/*") % (DATAX_HOME)
else:
CLASS_PATH= ("%s/lib/*:.") % (DATAX_HOME)
LOGBACK_FILE= ("%s/conf/logback.xml") % (DATAX_HOME)
DEFAULT_JVM="-Xms1g -Xmx1g -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=%s/log"% (DATAX_HOME)
DEFAULT_PROPERTY_CONF="-Dfile.encoding=UTF-8 -Dlogback.statusListenerClass=ch.qos.logback.core.status.NopStatusListener -Djava.security.egd=file:///dev/urandom -Ddatax.home=%s -Dlogback.configurationFile=%s"% (
DATAX_HOME, LOGBACK_FILE)
ENGINE_COMMAND="java -server ${jvm} %s -classpath %s ${params} com.alibaba.datax.core.Engine -mode ${mode} -jobid ${jobid} -job ${job}"% (
DEFAULT_PROPERTY_CONF, CLASS_PATH)
REMOTE_DEBUG_CONFIG="-Xdebug -Xrunjdwp:transport=dt_socket,server=y,address=9999"
RET_STATE= {
"KILL": 143,
"FAIL": -1,
"OK": 0,
"RUN": 1,
"RETRY": 2
}
defgetLocalIp():
try:
returnsocket.gethostbyname(socket.getfqdn(socket.gethostname()))
except:
return"Unknown"
defsuicide(signum, e):
globalchild_process
print>>sys.stderr, "[Error] DataX receive unexpected signal %d, starts to suicide."% (signum)
ifchild_process:
child_process.send_signal(signal.SIGQUIT)
time.sleep(1)
child_process.kill()
print>>sys.stderr, "DataX Process was killed ! you did ?"
sys.exit(RET_STATE["KILL"])
defregister_signal():
ifnotisWindows():
globalchild_process
signal.signal(2, suicide)
signal.signal(3, suicide)
signal.signal(15, suicide)
defgetOptionParser():
usage="usage: %prog [options] job-url-or-path"
parser=OptionParser(usage=usage)
prodEnvOptionGroup=OptionGroup(parser, "Product Env Options",
"Normal user use these options to set jvm parameters, job runtime mode etc. "
"Make sure these options can be used in Product Env.")
prodEnvOptionGroup.add_option("-j", "--jvm", metavar="<jvm parameters>", dest="jvmParameters", action="store",
default=DEFAULT_JVM, help="Set jvm parameters if necessary.")
prodEnvOptionGroup.add_option("--jobid", metavar="<job unique id>", dest="jobid", action="store", default="-1",
help="Set job unique id when running by Distribute/Local Mode.")
prodEnvOptionGroup.add_option("-m", "--mode", metavar="<job runtime mode>",
action="store", default="standalone",
help="Set job runtime mode such as: standalone, local, distribute. "
"Default mode is standalone.")
prodEnvOptionGroup.add_option("-p", "--params", metavar="<parameter used in job config>",
action="store", dest="params",
help='Set job parameter, eg: the source tableName you want to set it by command, '
'then you can use like this: -p"-DtableName=your-table-name", '
'if you have mutiple parameters: -p"-DtableName=your-table-name -DcolumnName=your-column-name".'
'Note: you should config in you job tableName with ${tableName}.')
prodEnvOptionGroup.add_option("-r", "--reader", metavar="<parameter used in view job config[reader] template>",
action="store", dest="reader",type="string",
help='View job config[reader] template, eg: mysqlreader,streamreader')
prodEnvOptionGroup.add_option("-w", "--writer", metavar="<parameter used in view job config[writer] template>",
action="store", dest="writer",type="string",
help='View job config[writer] template, eg: mysqlwriter,streamwriter')
parser.add_option_group(prodEnvOptionGroup)
devEnvOptionGroup=OptionGroup(parser, "Develop/Debug Options",
"Developer use these options to trace more details of DataX.")
devEnvOptionGroup.add_option("-d", "--debug", dest="remoteDebug", action="store_true",
help="Set to remote debug mode.")
devEnvOptionGroup.add_option("--loglevel", metavar="<log level>", dest="loglevel", action="store",
default="info", help="Set log level such as: debug, info, all etc.")
parser.add_option_group(devEnvOptionGroup)
returnparser
defgenerateJobConfigTemplate(reader, writer):
readerRef="Please refer to the %s document:\n https://github.com/alibaba/DataX/blob/master/%s/doc/%s.md \n"% (reader,reader,reader)
writerRef="Please refer to the %s document:\n https://github.com/alibaba/DataX/blob/master/%s/doc/%s.md \n "% (writer,writer,writer)
print(readerRef)
print(writerRef)
jobGuid='Please save the following configuration as a json file and use\n python {DATAX_HOME}/bin/datax.py {JSON_FILE_NAME}.json \nto run the job.\n'
print(jobGuid)
jobTemplate={
"job": {
"setting": {
"speed": {
"channel": ""
}
},
"content": [
{
"reader": {},
"writer": {}
}
]
}
}
readerTemplatePath="%s/plugin/reader/%s/plugin_job_template.json"% (DATAX_HOME,reader)
writerTemplatePath="%s/plugin/writer/%s/plugin_job_template.json"% (DATAX_HOME,writer)
try:
readerPar=readPluginTemplate(readerTemplatePath);
exceptExceptionase:
print("Read reader[%s] template error: can\'t find file %s"% (reader,readerTemplatePath))
try:
writerPar=readPluginTemplate(writerTemplatePath);
exceptExceptionase:
print("Read writer[%s] template error: : can\'t find file %s"% (writer,writerTemplatePath))
jobTemplate['job']['content'][0]['reader'] =readerPar;
jobTemplate['job']['content'][0]['writer'] =writerPar;
print(json.dumps(jobTemplate, indent=4, sort_keys=True))
defreadPluginTemplate(plugin):
withopen(plugin, 'r') asf:
returnjson.load(f)
defisUrl(path):
ifnotpath:
returnFalse
assert (isinstance(path, str))
m=re.match(r"^http[s]?://\S+\w*", path.lower())
ifm:
returnTrue
else:
returnFalse
defbuildStartCommand(options, args):
commandMap= {}
tempJVMCommand=DEFAULT_JVM
ifoptions.jvmParameters:
tempJVMCommand=tempJVMCommand+" "+options.jvmParameters
ifoptions.remoteDebug:
tempJVMCommand=tempJVMCommand+" "+REMOTE_DEBUG_CONFIG
print('local ip: ', getLocalIp())
ifoptions.loglevel:
tempJVMCommand=tempJVMCommand+" "+ ("-Dloglevel=%s"% (options.loglevel))
ifoptions.mode:
commandMap["mode"] =options.mode
# jobResource 可能是 URL,也可能是本地文件路径(相对,绝对)
jobResource=args[0]
ifnotisUrl(jobResource):
jobResource=os.path.abspath(jobResource)
ifjobResource.lower().startswith("file://"):
jobResource=jobResource[len("file://"):]
jobParams= ("-Dlog.file.name=%s") % (jobResource[-20:].replace('/', '_').replace('.', '_'))
ifoptions.params:
jobParams=jobParams+" "+options.params
ifoptions.jobid:
commandMap["jobid"] =options.jobid
commandMap["jvm"] =tempJVMCommand
commandMap["params"] =jobParams
commandMap["job"] =jobResource
returnTemplate(ENGINE_COMMAND).substitute(**commandMap)
defprintCopyright():
print('''
DataX (%s), From Alibaba !
Copyright (C) 2010-2017, Alibaba Group. All Rights Reserved.
'''%DATAX_VERSION)
sys.stdout.flush()
if__name__=="__main__":
printCopyright()
parser=getOptionParser()
options, args=parser.parse_args(sys.argv[1:])
ifoptions.readerisnotNoneandoptions.writerisnotNone:
generateJobConfigTemplate(options.reader,options.writer)
sys.exit(RET_STATE['OK'])
iflen(args) !=1:
parser.print_help()
sys.exit(RET_STATE['FAIL'])
startCommand=buildStartCommand(options, args)
# print startCommand
child_process=subprocess.Popen(startCommand, shell=True)
register_signal()
(stdout, stderr) =child_process.communicate()
sys.exit(child_process.returncode)