菜单

SDH importer

1. 应用场景

客户的数据被写成了神策的 JSON 格式(数据格式)被放到了神策的 HDFS 文件 / 本地文件上。希望使用工具将数据导入神策的实体中,则使用本工具写入数据。

2. SDH importer 使用步骤 

  1. 切换到 sa_cluster 账户
    1. sudo su - sa_cluster
  2. 准备数据文件
    1. 文件内容为每行一个符合 数据格式 的 json
    2. 文件可以使用 gzip 压缩,压缩后的文件名需要以 ".gz" 为后缀。例如,原始 json 文件的文件名为 "data1.json",使用 gzip 压缩后,文件名需要改为 "data1.json.gz"。
  3. 创建 HDFS 目录,将数据文件上传至 HDFS。一次导入的多个数据文件,需要放在同一个 HDFS 目录下。

    $ hdfs dfs -put ${path_to_local_file} ${path_to_hdfs_dir}
  4. 执行导入任务
    1. 提交数据导入任务

      # 导入数据
      horizonadmin inflow importer_v2 submit_task --path hdfs:///sa/foo/bar/ --project_name default
    2. 在命令输出中,可以获取任务的 task_id
  5. 获取已提交任务的执行状态
    1. # 获取任务
      horizonadmin inflow importer_v2 get_task --task_id 6
  6. 取消任务 
    1. # 取消任务
      horizonadmin inflow importer_v2 cancel_task --task_id 6
    2. 注意:若任务已经执行到导入后期,可能数据已经入库
  7. 重试任务
    1.  
      # 重试任务
      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. 要导入一大批数据,导入性能如何调优?

  1. 确认可用资源
    1. 批数据导入操作,主要占用的是 Yarn 资源,包括 Yarn 可使用的 CPU 及内存,首先确保 Yarn 可使用的资源充足。
  2. 调整任务的并发度和内存
    1. 通过 --parallelism 调整执行并行度,尽可能提高资源使用率
  3. 提前创建元数据,使用 --data_to_parquet_fast_mode 参数
    1. 若能提前创建元数据,确认事件属性和用户属性,则可以在提交导入任务时使用 --data_to_parquet_fast_mode,跳过属性检查步骤,提升导入效率。
  4. 通过任务执行状态,判断执行瓶颈
    1. 使用 horizonadmin inflow importer_v2 get_task 命令获取已经执行的任务的状态,可以获取任务执行中不同阶段的 QPS,以此找神策专家判断是否有优化空间。
上一个
数据导入常见问题
下一个
数据导出
最近修改: 2026-08-17