【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文基于 Apache Beam 仓库中的CdapServiceNowToTxt批处理示例,讲解如何借助 CDAP 的 ServiceNow 插件(通过CdapIO集成),从 ServiceNow 实例批量拉取 JSON 格式数据并写入本地.txt文件。读完本文,你将掌握:如何配置 Gradle 执行任务、如何填写 10 个核心管道参数(认证信息、查询模式、值类型、输出路径等)、如何理解整条管道的源码级执行链路,以及如何切换到不同 Runner 运行该管道。
该示例位于 examples/java/cdap/servicenow/src/main/java/org/apache/beam/examples/complete/cdap/servicenow,属于 examples 下 CDAP 插件示例集合的一部分(cdap 示例总览),同目录下还有 Salesforce、Hubspot、Zendesk 等同类示例。
一、示例概览:从 ServiceNow 到 .txt 的批处理管道
CdapServiceNowToTxt是一条纯批处理管道,目标非常明确:以 JSON 格式从 CDAP ServiceNow 源读取数据,把结果记录写入 .txt 文件。ServiceNow 相关参数与输出文件路径均由用户在模板参数中指定。
1.1 类定位与入口
主类为 CdapServiceNowToTxt.java,其main方法负责两件事:
public static void main(String[] args) { CdapServiceNowOptions options = PipelineOptionsFactory.fromArgs(args).withValidation().as(CdapServiceNowOptions.class); // Create the pipeline Pipeline pipeline = Pipeline.create(options); run(pipeline, options); }- 通过
PipelineOptionsFactory.fromArgs(args).withValidation()解析命令行参数,并开启参数校验——所有标了@Validation.Required的参数若缺失会直接报错; - 将参数转换为
CdapServiceNowOptions类型后创建Pipeline并调用run(pipeline, options)执行。
1.2 三步式执行流程
run方法中的注释清晰地概括了管道核心步骤:
/* * Steps: * 1) Read messages in from Cdap ServiceNow * 2) Extract values only * 3) Write successful records to .txt file */对应到代码实现,是一条由三个 Transform 串起来的流水线:
pipeline .apply("readFromCdapServiceNow", FormatInputTransform.readFromCdapServiceNow(paramsMap)) .setCoder(KvCoder.of( NullableCoder.of(WritableCoder.of(NullWritable.class)), SerializableCoder.of(StructuredRecord.class))) .apply(MapValues.into(TypeDescriptors.strings()) .via(StructuredRecordUtils::structuredRecordToString)) .setCoder(KvCoder.of( NullableCoder.of(WritableCoder.of(NullWritable.class)), StringUtf8Coder.of())) .apply(Values.create()) .apply("writeToTxt", TextIO.write().to(options.getOutputTxtFilePathPrefix()));各环节说明:
| 步骤 | Transform | 作用 |
|---|---|---|
| 1 | FormatInputTransform.readFromCdapServiceNow(paramsMap) | 通过CdapIO.read()加载 ServiceNow CDAP 插件,读取数据,产出KV<NullWritable, StructuredRecord> |
| 2 | MapValues+StructuredRecordUtils::structuredRecordToString | 仅提取 Value 部分,把StructuredRecord序列化为 JSON 字符串 |
| 3 | Values.create()+TextIO.write() | 丢弃 Key,把每条 JSON 记录写入以指定前缀命名的.txt文件 |
其中第 2 步用到的structuredRecordToString定义在 StructuredRecordUtils.java,内部使用 Gson 将StructuredRecord转为 JSON;若记录为null则输出"{}",保证下游写出时不会出现空记录崩溃。
1.3 关键点:CDAP 插件驱动的数据源
FormatInputTransform.readFromCdapServiceNow(定义在 FormatInputTransform.java)是数据接入的核心:
final PluginConfig pluginConfig = new ConfigWrapper<>(ServiceNowSourceConfig.class).withParams(pluginConfigParams).build(); checkStateNotNull(pluginConfig, "Plugin config can't be null."); return CdapIO.<NullWritable, StructuredRecord>read() .withCdapPluginClass(ServiceNowSource.class) .withPluginConfig(pluginConfig) .withKeyClass(NullWritable.class) .withValueClass(StructuredRecord.class);这段代码揭示了底层机制:
- 用
ConfigWrapper<ServiceNowSourceConfig>把参数 Map 包装成 CDAP 插件配置对象; - 通过
CdapIO.read()指定插件类ServiceNowSource(来自io.cdap.plugin.servicenow.source包),即以 CDAP 生态中的 ServiceNow 批处理 Source 插件作为数据源; - Key 类型为 Hadoop 的
NullWritable,Value 类型为 CDAP 的StructuredRecord,这与管道中设置的 Coder 一一对应。
也就是说,该示例本身不直接实现 ServiceNow REST 调用逻辑,而是复用 CDAP 数据集成生态中成熟的 ServiceNow 插件,这正是 Apache BeamCdapIO桥接层(org.apache.beam.sdk.io.cdap)的典型用法。
二、Gradle 准备:创建执行任务
在运行示例前,需要在项目的build.gradle中声明一个 JavaExec 任务,用于动态指定主类与命令行参数:
task executeCdapServiceNow (type:JavaExec) { mainClass = System.getProperty("mainClass") classpath = sourceSets.main.runtimeClasspath systemProperties System.getProperties() args System.getProperty("exec.args", "").split() }该任务的核心设计思路:
| 配置项 | 含义 |
|---|---|
mainClass = System.getProperty("mainClass") | 从系统属性读取主类全限定名,方便在命令行用-DmainClass=...动态指定 |
classpath = sourceSets.main.runtimeClasspath | 使用主源码集的运行时 classpath,保证 CDAP 插件、Beam SDK 等依赖可用 |
systemProperties System.getProperties() | 把 JVM 系统属性透传给子进程 |
args System.getProperty("exec.args", "").split() | 把-Dexec.args传入的参数字符串按空格拆分成参数数组;未提供时为空 |
三、运行 CdapServiceNowToTxt 管道
3.1 基本运行命令
Gradle 任务executeCdapServiceNow支持通过如下命令运行管道:
gradle clean executeCdapServiceNow -DmainClass=org.apache.beam.examples.complete.cdap.servicenow.CdapServiceNowToTxt \ -Dexec.args="--<argument>=<value> --<argument>=<value>"要点说明:
clean用于清理上次构建产物,避免旧 class 干扰;-DmainClass指定入口类org.apache.beam.examples.complete.cdap.servicenow.CdapServiceNowToTxt;-Dexec.args内以--key=value形式传递所有管道参数,参数之间用空格分隔。
3.2 完整参数示例
执行管道时需按如下格式指定参数:
--clientId=your-client-id \ --clientSecret=your-client-secret \ --user=your-user \ --password=your-password \ --restApiEndpoint=your-endpoint \ --queryMode=Table \ --tableName=your-table \ --valueType=Actual \ --referenceName=your-reference-name \ --outputTxtFilePathPrefix=your-path-to-output-folder-with-filename-prefix3.3 切换 Runner
默认情况下管道使用 DirectRunner 在本地执行。如需更换执行引擎,追加:
--runner=YOUR_SELECTED_RUNNER例如在本地验证时用--runner=DirectRunner(默认),在分布式环境可换成DataflowRunner、FlinkRunner、SparkRunner等 Apache Beam 支持的 Runner。
四、参数详解:9 个必填参数 + 1 个参考名
所有参数定义集中在 CdapServiceNowOptions.java,该接口继承自 BaseCdapOptions.java,后者统一提供了referenceName参数。下面逐一解析:
4.1 认证类参数(4 个)
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
clientId | String | ✅ | ServiceNow 实例的 Client ID(OAuth 客户端标识) |
clientSecret | String | ✅ | ServiceNow 实例的 Client Secret(OAuth 客户端密钥) |
user | String | ✅ | ServiceNow 实例的用户名 |
password | String | ✅ | ServiceNow 实例的密码 |
源码中以@Validation.Required标注,运行时若缺失会触发校验失败。安全提示:明文在命令行中传递密码存在泄露风险,生产环境建议通过受保护的环境变量或密钥管理机制注入。
4.2 连接与查询类参数(3 个)
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
restApiEndpoint | String | ✅ | ServiceNow 实例的 REST API 端点,例如https://instance.service-now.com |
queryMode | String | ✅ | 查询模式,二选一:Reporting—— 选择应用后拉取该应用下所有表的数据;Table—— 直接指定表名拉取数据 |
tableName | String | ✅ | 要读取数据的 ServiceNow 表名;当queryMode=Reporting时该值会被忽略 |
4.3 取值与输出类参数(3 个)
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
valueType | String | ✅ | 返回值的类型:Actual—— 返回表中实际存储值(默认);Display—— 返回表的显示值 |
referenceName | String | ✅ | 参考名(来自 BaseCdapOptions),CDAP 插件通用配置,用于标识该数据源 |
outputTxtFilePathPrefix | String | ✅ | 输出文件夹路径 + 文件名前缀;写出时会生成{prefix}-###形式的一组 .txt 文件 |
outputTxtFilePathPrefix的###分片命名规则来自TextIO.write()的标准行为,多个 Worker 并行写出时会得到{prefix}-00000-of-00001、{prefix}-00000-of-00002等文件。
五、源码级深入:参数如何流向 CDAP 插件
5.1 参数 → 插件配置 Map 的转换
PluginConfigOptionsConverter.java 负责把 Pipeline Options 转成 CDAP 插件认识的参数 Map:
return ImmutableMap.<String, Object>builder() .put(ServiceNowConstants.PROPERTY_CLIENT_ID, options.getClientId()) .put(ServiceNowConstants.PROPERTY_CLIENT_SECRET, options.getClientSecret()) .put(ServiceNowConstants.PROPERTY_USER, options.getUser()) .put(ServiceNowConstants.PROPERTY_PASSWORD, options.getPassword()) .put(ServiceNowConstants.PROPERTY_API_ENDPOINT, options.getRestApiEndpoint()) .put(ServiceNowConstants.PROPERTY_QUERY_MODE, options.getQueryMode()) .put(ServiceNowConstants.PROPERTY_TABLE_NAME, options.getTableName()) .put(ServiceNowConstants.PROPERTY_VALUE_TYPE, options.getValueType()) .put(Constants.Reference.REFERENCE_NAME, options.getReferenceName()) .build();可见每个命令行参数都被映射为ServiceNowConstants中的插件属性常量(如PROPERTY_CLIENT_ID、PROPERTY_QUERY_MODE),referenceName则映射为 CDAP 通用常量Constants.Reference.REFERENCE_NAME。
5.2 完整调用链
命令行 --clientId=... --queryMode=Table ... │ PipelineOptionsFactory 解析 ▼ CdapServiceNowOptions │ PluginConfigOptionsConverter.serviceNowOptionsToParamsMap() ▼ Map<String, Object> paramsMap │ FormatInputTransform.readFromCdapServiceNow(paramsMap) ▼ ConfigWrapper<ServiceNowSourceConfig> → CdapIO.read() │ 指定 ServiceNowSource 插件 ▼ KV<NullWritable, StructuredRecord> 数据流 │ MapValues → StructuredRecordUtils.structuredRecordToString ▼ JSON 字符串 → TextIO.write() → {prefix}-###.txt这条链路完整展示了 Beam 示例如何"寄生"在 CDAP 插件体系之上:Beam 负责参数解析、Coder 设置与写出,CDAP 插件负责与 ServiceNow REST API 交互。
六、输出格式与验证
- 输出内容:每行一条 JSON 格式的记录,对应 ServiceNow 表的一行数据(由
structuredRecordToString通过 Gson 序列化生成)。 - 输出文件:以
--outputTxtFilePathPrefix指定的前缀生成的多个分片文件,形如{prefix}-00000-of-00001。
验证方式:管道运行结束后,检查输出目录下生成的.txt文件,确认每行 JSON 记录中的字段与 ServiceNow 表中查询到的数据一致;若valueType=Display,字段值应为表定义的显示值而非存储值。
七、注意事项与限制
- 参数校验严格:9 个必填参数(
clientId、clientSecret、user、password、restApiEndpoint、queryMode、tableName、valueType、outputTxtFilePathPrefix)加上referenceName共 10 项,全部标注@Validation.Required,缺失任一参数管道启动即失败。 - 依赖 CDAP 插件生态:本示例依赖
io.cdap.plugin.servicenow相关 JAR,构建时需确保该依赖在 classpath 中(可通过classpath = sourceSets.main.runtimeClasspath自动带入)。 queryMode与tableName联动:使用Reporting模式时tableName被忽略,选择应用后拉取该应用下所有表的数据;使用Table模式时则必须指定具体表名。- 凭据安全:认证信息通过命令行明文传递,适合本地演示与测试;生产环境应结合安全配置管理。
- Runner 兼容性:示例本身为标准 Beam 批处理管道(
TextIO+MapValues),理论上可运行在任意支持批处理的 Beam Runner 上;本文以本地 DirectRunner 为例说明。
八、进一步探索
- 阅读其他 CDAP 插件示例:可对比 Salesforce、Hubspot、Zendesk 等目录下的同构实现,加深对
CdapIO通用接入模式的理解。 - 深入
CdapIO底层:其桥接实现位于org.apache.beam.sdk.io.cdap包(Java SDK 的 core 相关源码),可进一步研究CdapIO.Read如何基于 CDAP 插件构建 Beam 数据源。 - 如果需求是从 ServiceNow 读取数据后再做转换、聚合或写入其他存储(如 BigQuery、GCS),只需在上述三步流水线中插入对应 Beam Transform 即可,本示例是最简化的可扩展骨架。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Java Kata 实战:用 TextIO.read() 从文本文件读取 PCollection
Apache Beam Java Kata 实战:用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官
大数据批处理流处理数据工程Apache Beam CdapIO 实战:构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例
Apache Beam CdapIO 实战:构建从 ServiceNow 批量拉取数据并落盘 TXT 的管道示例 本文基于 Apache Beam 仓库中 ex
大数据批处理流处理数据工程Apache Beam Kotlin Kata 实战:使用 TextIO 从文本文件读取 PCollection
Apache Beam Kotlin Kata 实战:使用 TextIO 从文本文件读取 PCollection 创建 Beam 管道时,最常见的第一步就是从文
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考