1. 应用场景
客户的数据被写成了神策的 JSON 格式(数据格式)被放到了神策的 HDFS 文件 / 本地文件上。希望使用工具将数据导入神策的实体中,则使用本工具写入数据。
2. SDH importer 使用步骤
- 切换到 sa_cluster 账户
-
sudo su - sa_cluster
-
- 准备数据文件
- 文件内容为每行一个符合 数据格式 的 json
- 文件可以使用 gzip 压缩,压缩后的文件名需要以 ".gz" 为后缀。例如,原始 json 文件的文件名为 "data1.json",使用 gzip 压缩后,文件名需要改为 "data1.json.gz"。
-
创建 HDFS 目录,将数据文件上传至 HDFS。一次导入的多个数据文件,需要放在同一个 HDFS 目录下。
$ hdfs dfs -put ${path_to_local_file} ${path_to_hdfs_dir} -
- 执行导入任务
-
提交数据导入任务
# 导入数据 horizonadmin inflow importer_v2 submit_task --path hdfs:///sa/foo/bar/ --project_name default - 在命令输出中,可以获取任务的 task_id
-
- 获取已提交任务的执行状态
-
# 获取任务 horizonadmin inflow importer_v2 get_task --task_id 6
-
- 取消任务
-
# 取消任务 horizonadmin inflow importer_v2 cancel_task --task_id 6 - 注意:若任务已经执行到导入后期,可能数据已经入库
-
- 重试任务
-
# 重试任务 horizonadmin inflow importer_v2 retry_task --task_id 6
-
兼容性说明
- 与 BatchImporter 的兼容性
- 由于老架构的 BatchImporter 在使用上有诸多问题,新版本中予以优化,在使用上也有一定的差别,详见:
- 与 integratoradmin importer 的兼容性
- “horizonadmin inflow importer_v2” 命令,与 "integratoradmin importer" 命令完全兼容,升级不受影响。若升级前使用了"integratoradmin importer" 工具导入数据,升级后依然可用。
3. 命令参数详情
3.1. 提交任务命令参数详情
使用 "horizonadmin inflow importer_v2 submit_task" 命令来提交导入任务,并等待任务执行完成
|
参数 |
是否必选 |
默认值 |
描述 |
样例 |
|---|---|---|---|---|
|
--path |
必选其一 |
- |
批量导入 hdfs 数据的文件夹路径 |
--path hdfs:///sa/runtime/data_dir |
|
--local_path |
- |
批量导入本地数据的路径 |
--local_path /home/sa_cluster/1.txt |
|
|
--name |
否 |
UUID |
任务名,需要项目内唯一 |
--name test_event_20230415_01 |
|
--project_name |
必选其一 |
- |
项目名 |
--project_name production |
|
--project_id |
- |
项目 id |
--project_id 2 |
|
|
--expired_record_filter_after_hour |
否 |
17520(2 * 365 * 24) |
数据有效时间,向前偏移,单位小时 |
--expired_record_filter_after_hour 35040 |
|
--expired_record_filter_before_hour |
否 |
1 |
数据有效时间,向后偏移,单位小时 |
--expired_record_filter_before_hour 24 |
|
--enable_first_time_processor |
否 |
true |
开启首次导入处理 |
--enable_first_time_processor true |
|
--data_to_parquet_fast_mode |
否 |
- |
数据文件直接转化为 parquet 格式 使用场景:在 提前声明好所有事件属性及用户属性 后,使用此参数,可跳过元数据检查的步骤,加快导入速度。 ⚠注:
|
--data_to_parquet_fast_mode |
|
--retry_times |
否 |
0 |
重试次数,0 表示不重试 接入初期不建议开重试,会影响问题排查 |
--retry_times 0 |
|
--created_by |
否 |
CONSOLE |
来源,标记任务来源,便于溯源 |
--created_by SHOPPING_MALL |
|
--parallelism |
否 | 3 | 执行并行度 |
--parallelism 10 |
|
--custom_processor_jar_path |
否 |
- |
【高级参数】预处理 jar 包路径 |
--custom_processor_jar_path /home/sa_cluster/1.jar |
|
--custom_processor |
否 |
- |
【高级参数】预处理 jar 的类名 |
--custom_processor com.sensorsdata.analytics.csvtojson.CSVToJSON |
|
--custom_processor_config |
否 |
- |
【高级参数】预处理 jar 的配置 |
--custom_processor_config abc |
|
--flink_conf: |
否 |
集群默认值: jobmanager 内存 1024M taskmanager 内存 2048M 并行度 3 单机默认值: 并行度 1 |
【高级参数】flink config 的动态配置 |
--flink_conf:jobmanager.memory.process.size=2048m --flink_conf:taskmanager.memory.process.size=4096m --flink_conf:parallelism.default=3 |
|
--enabled_steps |
否 |
- |
【高级参数】允许执行的步骤 |
--enabled_steps BATCH_PROCESS,BATCH_LOAD_TO_PARQUET,BATCH_LOAD_TO_SDW,REMAPPING_SEND,SUPPLY_DATA |
|
--ytm |
否 | - |
【高级参数】cluster 模式, Yarn TaskManager 占用内存,单位为 MB |
--ytm 4096 |
|
--run_async |
否 |
- |
【高级参数】异步执行 |
--run_async |
3.2. 查询任务参数详情
使用 "horizonadmin inflow importer_v2 get_task" 命令查询任务
|
参数 |
是否必选 |
默认值 |
描述 |
样例 |
|---|---|---|---|---|
|
--task_id |
task_id / task_name + project 必传 |
- |
--task_id 6 |
|
|
--task_name |
- |
--task_name test_name |
||
|
--project |
- |
项目英文名 |
--project default |
|
|
--simple |
false |
展示简洁信息 |
--simple |
3.3. 重试任务参数详情
使用 "horizonadmin inflow importer_v2 retry_task" 命令重试任务
|
参数 |
是否必选 |
默认值 |
描述 |
样例 |
|---|---|---|---|---|
| --task_id | task_id / task_name + project 必传 |
- | --task_id 6 | |
| --task_name | - | --task_name test_name | ||
| --project | - | 项目英文名 | --project default | |
| --enabled_steps | false | 允许执行的步骤 | --enabled_steps BATCH_PROCESS,BATCH_LOAD_TO_PARQUET,BATCH_LOAD_TO_SDW,REMAPPING_SEND,SUPPLY_DATA |
3.4. 取消任务参数详情
使用 "horizonadmin inflow importer_v2 cancel_task" 命令取消任务
|
参数 |
是否必选 |
默认值 |
描述 |
样例 |
|---|---|---|---|---|
|
参数 |
是否必选 |
默认值 |
描述 |
样例 |
|
--task_id |
task_id / task_name + project 必传 |
- |
--task_id 6 |
|
|
--task_name |
- |
--task_name test_name |
||
|
--project |
- |
项目英文名 |
--project default |
4. FAQ
4.1. 如果数据导错了,想删除重导,要如何操作?
如果是偶然一次导入失误需要删除,可以联系神策值班处理。
若是频繁有导入删除的操作需求,最佳的方案是:在导入数据中,自定义一个字段,标识“数据导入批次”。当发现数据错误时,可使用 【SDH tools】数据删除工具(高危) 进行数据删除。
4.2. 要导入一大批数据,导入性能如何调优?
- 确认可用资源
- 批数据导入操作,主要占用的是 Yarn 资源,包括 Yarn 可使用的 CPU 及内存,首先确保 Yarn 可使用的资源充足。
- 调整任务的并发度和内存
- 通过 --parallelism 调整执行并行度,尽可能提高资源使用率
- 提前创建元数据,使用 --data_to_parquet_fast_mode 参数
- 若能提前创建元数据,确认事件属性和用户属性,则可以在提交导入任务时使用 --data_to_parquet_fast_mode,跳过属性检查步骤,提升导入效率。
- 通过任务执行状态,判断执行瓶颈
- 使用 horizonadmin inflow importer_v2 get_task 命令获取已经执行的任务的状态,可以获取任务执行中不同阶段的 QPS,以此找神策专家判断是否有优化空间。