
SeaTunnel SensorsData Sink Connector 实战指南基于官方 SDK 的用户事件、用户档案与物品数据上报【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文是一份基于 Apache SeaTunnel 官方文档与仓库源码的深度技术指南全面讲解 SensorsData Sink Connector 的接入原理、全部配置参数、条件校验规则与 6 类典型使用场景。读完本文你将能够基于 SeaTunnel 的 HOCON 配置把任意上游数据流以用户事件events、用户档案users/profile、用户明细details和物品记录items四种形态写入神策数据SensorsData服务端并掌握console本地调试、${field}动态事件名、current_time()处理时间、time_free历史数据模式、null_as_profile_unset档案字段清除等高级用法。一、Connector 定位与能力总览SensorsData Sink 是 SeaTunnel Connector V2 体系中的一个写端插件位于 seatunnel-connectors-v2/connector-sensorsdata 模块。它不直接拼装 HTTP 请求而是复用神策官方 Java SDKcom.sensorsdata.analytics.javasdk将 SeaTunnel 内部的行数据SeaTunnelRow转换为神策的数据模型后交给 SDK 发送。支持的引擎SensorsData Sink 基于引擎无关的 SeaTunnel Connector API 实现因此天然支持以下执行引擎SeaTunnel Zeta推荐FlinkSpark能力边界来自官方文档标注exactly-once不支持。从源码看SensorsDataSDKWriter 仅在prepareCommit()中执行sa.flush()并返回空的SensorsDataCommitInfo没有两阶段提交语义因此该 Sink 属于at-least-once级别cdc增量变更捕获不支持Sink 不会根据行 kindINSERT/UPDATE/DELETE区分写入策略。支持写入的数据类型根据 SensorsDataOptions 与 TypeUtil 实现SensorsData Sink 可处理以下四类实体数据用户事件USER_EVENT——调用 SDK 的sa.track(...)对应神策埋点事件用户档案USER——调用sa.profileSet(...)并在开启null_as_profile_unset时额外调用sa.profileUnset(...)用户明细USER_DETAIL——调用sa.detailSet(...)物品记录SPECIAL_ITEM——调用sa.itemSet(...)兼容非 SDH 架构的双主键物品表。说明SensorsDataRecordType枚举中还定义了ITEM、ITEM_EVENT、ITEM_DETAIL三个类型但源码注释明确标注Not Implemented Yet尚未实现从SensorsDataRecordBuilder的分支逻辑看entity_name items目前只会走到SPECIAL_ITEM分支。二、Sink 选项完整参考下表为官方文档定义的完整参数表已结合 SensorsDataOptions 与 SensorsDataSDKSinkOptions 源码校验NameTypeRequiredDefaultDescriptionserver_urlstringYes-神策数据接收地址格式https://host:8106/sa?projectdefaultbulk_sizeintNo50SensorsData SDK 缓存的触发 flush 阈值max_cache_row_sizeintNo0超过该值立即触发 flush0表示跟随bulk_sizeconsumerstringNobatchconsumer 类型batch发送到神策服务端console仅本地打印转换后的记录entity_namestringYesusers实体名支持users与itemsrecord_typestringYesusers记录类型常用events、details、items用户档案直接用usersschemastringConditionalusers数据表schema名称用户记录entity_name users必填distinct_id_columnstringConditional-作为神策 distinct_id 的输入列用户记录必填identity_fieldsarrayConditional-身份标识映射如$identity_login_id、$identity_distinct_id用户记录必填property_fieldsarrayConditional-属性映射与目标数据类型用户记录必填event_namestringConditional-事件名或${field_name}表达式record_type events时必填time_columnstringConditional-事件时间列record_type events时必填填current_time()使用处理时间detail_id_columnstringConditional-明细 ID 列record_type details时必填item_id_columnstringConditional-物品 ID 列record_type items时必填item_type_columnstringConditional-物品类型列record_type items时必填time_freebooleanNofalse是否开启神策无时区历史数据模式skip_error_recordbooleanNofalse转换失败的记录是否跳过true跳过false抛异常终止instant_eventsarrayNo[]需要标记为即时事件的事件名列表distinct_id_by_identitiesbooleanNofalse当distinct_id_column取值为空时自动从identity_fields中填充 distinct_idnull_as_profile_unsetbooleanNofalse将档案属性中的 null 值转换为 profile_unset 操作common-optionsconfigNo-Sink 通用选项见 Sink Common Options条件必填规则OptionRule 校验从 SensorsDataBaseOptionRules 可以看到参数之间的依赖是通过OptionRule.conditional(...)声明式校验的这解释了上表中 Conditional 的含义entity_name与record_type为始终必填当entity_name users时schema、distinct_id_column、identity_fields、property_fields成为必填当record_type events时time_column、event_name成为必填当record_type details时detail_id_column成为必填当record_type items时item_id_column、item_type_column成为必填time_free为可选。此外SensorsDataSDKSinkFactory.optionRule() 还会追加bulk_size、max_cache_row_size、skip_error_record、instant_events为可选参数。类型转换与必填字段的底层强校验值得注意除了配置层校验RowAccessor 在 Writer 初始化阶段还会做列级运行时校验所有配置的列名distinct_id_column、time_column、event_name中的${field}、detail_id_column、item_id_column、item_type_column以及property_fields/identity_fields中的source都必须在上游 schema 中存在否则抛UNKNOWN_SOURCE_FIELDschema、detailId、itemId、itemType在对应的记录类型下取值为空时会抛MISSING_NECESSARY_FIELD。这意味着列名拼写错误会在作业启动阶段就被拦截而不是等到数据写出时才发现这是排障时的第一检查点。三、参数深度解读server_url [string]神策数据接收地址格式固定为https://${host}:8106/sa?project${project}其中project为神策项目标识。示例https://10.1.136.63:8106/sa?projectdefault。该参数无默认值为必填项。bulk_size [int] 与 max_cache_row_size [int]这两个参数共同控制 SDK 内存缓存的刷写flush行为bulk_size默认50内存缓存队列达到该条数阈值时触发批量发送max_cache_row_size默认0最大缓存刷新条数一旦超过立即触发 flush。0表示跟随bulk_size的语义。从 SensorsDataSDKWriter 的构造逻辑可以看到它们被直接透传给神策 SDK 的BatchConsumer(serverUrl, bulkSize, maxCacheRowSize, false, 3, instantEvents)。另外Writer 的prepareCommit()中会显式调用sa.flush()即每次 checkpoint 也会强制刷写一次缓存与阈值刷写共同构成发送触发条件。consumer [string]默认值batch通过BatchConsumer批量上报到server_url设为consoleWriter 会改用ConsoleConsumer(new PrintWriter(System.out))将转换后的神策记录打印到本地控制台不会发送到服务端非常适合在不搭环境的情况下验证字段映射是否正确。源码中 consumer 比较是忽略大小写的CONSUMER_TYPE_CONSOLE.equalsIgnoreCase(...)。entity_name / record_type / schema [string]这三个参数共同定位神策实体数据模型中的目标表entity_name实体名取值users默认或itemsrecord_type记录类型取值users/events/details/itemsschema数据表名users实体下必填。从 SensorsDataRecordBuilder 的 Builder 构造逻辑可以还原出entity_name×record_type的组合规则entity_namerecord_type内部记录类型对应 SDK 方法usersusersUSERprofileSet/profileUnsetuserseventsUSER_EVENTtrackusersdetailsUSER_DETAILdetailSetitemsitemsSPECIAL_ITEMitemSet其他组合--抛UNSUPPORTED_RECORD_TYPE注意record_type在配置层默认值是users但如果你配置了entity_name items而record_type仍保持默认的users会触发不支持异常文档与示例中record_type items时entity_name应显式配置为items文档示例中省略了entity_name但源码要求entity_name必填建议显式写出。distinct_id_column [string]用户实体users的 distinct_id 来源列。从 RowAccessor.getDistinctId 的实现看该列值会统一转换为 STRING 类型后作为神策的distinct_id。identity_fields [array]用户实体的身份标识映射列表每个元素形如{ source ${来源列}, target ${神策身份字段} }。常用目标身份字段包括$identity_login_id——登录 ID$identity_anonymous_id——匿名 ID$identity_distinct_id——去重 ID$identity_email、$identity_phone等自定义身份。从 RowAccessor.getUserIdentities 的实现细节看$identity_login_id被转换为 STRING其余身份字段被转换为 LIST这与神策一个用户可以拥有多个匿名 ID/邮箱/手机号的身份模型一致值为 null 或空白字符串的身份会被自动跳过。property_fields [array] 与支持的类型属性映射列表每个元素形如{ target ${目标属性名}, source ${来源列}, type ${数据类型} }。官方文档列出的支持类型为BOOLEANDECIMALINTBIGINTFLOATDOUBLENUMBERSTRINGDATETIMESTAMPLISTLIST_COMMALIST_SEMICOLON从 TypeUtil.toTargetType 可以补充一些底层转换细节帮助你在字段类型不匹配时快速定位问题BOOLEAN接受 Boolean 原值、数值0/0.0/0L视为 false其余为 true、以及忽略大小写的字符串trueINT/BIGINT/FLOAT/DOUBLE/NUMBER统一走toNumber支持 Number 原值、可解析的数字字符串、Booleantrue→1false→0DECIMAL数字字符串用BigDecimal解析Boolean 转 1/0STRING直接toString()byte[]会先转为字符串TIMESTAMPDate/LocalDate/LocalDateTime/Number直接换算成毫秒时间戳字符串会按内部格式数组依次尝试解析yyyy-MM-dd HH:mm:ss.SSS、yyyy-MM-dd HH:mm:ss、yyyy-MM-dd HH:mm、yyyy-MM-dd、yyyyMMdd_HHmmss、yyyyMMdd解析失败时保留原始值并打印 warn 日志DATE字符串按同样格式集解析后转为java.util.DateLIST / LIST_COMMA / LIST_SEMICOLON仅接受字符串类型分别以换行符\n、逗号,、分号;切分非字符串输入会抛DATA_TYPE_CAST_FIELD异常。注意TIME_COLUMN对应神策系统属性$time从getProperties的逻辑看配置了时间列时该列会被转换为 DATE 类型并写入$time系统属性配置current_time()时则写入当前时间new Date()。event_name [string]——静态与动态两种格式支持两种写法静态事件名直接填写事件名称例如event_name $AppStart动态事件名使用${字段名}表达式事件名取上游数据中该列的值。例如上游数据如下nameprop1prop2Purchase16data-example1Order23data-example2若配置event_name ${name}则第一行事件名为Purchase第二行事件名为Order。从 RowAccessor.initEventNameConfig 可以看到其实现配置值用正则\$\{(.*?)\}匹配命中则解析为动态列索引未命中则视为静态事件名。因此${...} 必须在 event_name 中成对出现且列名必须存在于上游 schema。time_column [string] 与 time_free [boolean]time_column事件时间列。除了上游时间字段还支持特殊值current_time()表示使用处理时间即 SeaTunnel 处理该行数据的当前时刻见RowAccessor中的CURRENT_TIME_KEY与getProperties逻辑time_free开启神策无时区历史数据模式。从 UserRecordBase 的addTimeFree看仅当记录类型为 track 事件且开启time_free时序列化 JSON 中会追加time_free: true字段。detail_id_column / item_id_column / item_type_column [string]detail_id_column明细记录的唯一标识列record_type details时必填item_id_column物品 ID 列record_type items时必填item_type_column物品类型列record_type items时必填。三者在运行时都会经过非空强校验取值为空白时抛MISSING_NECESSARY_FIELD。skip_error_record [boolean]默认false。从 SensorsDataSDKWriter.write 的实现看每条记录的构建/发送都包在 try-catch 中转换或发送异常时先打印 error 日志含 tableId、rowKind 与逐字段值skip_error_record true时仅记录日志、跳过该行作业继续false时抛出SEND_RECORD_FAILED异常终止任务。instant_events [array]事件名列表列表中的事件会被标记为即时事件。从 Writer 构造逻辑可见该列表被透传给神策 SDK 的BatchConsumer构造参数第五个参数由 SDK 层实现即时上报语义。distinct_id_by_identities [boolean]默认false。开启后当distinct_id_column取值为空时自动从identity_fields中填充 distinct_id。从 RowAccessor.getDistinctId 的备选逻辑看填充顺序为distinct_id→$identity_login_id→$identity_anonymous_id→$identity_distinct_id→ 其他身份字段以字段名值拼接保证神策收到非空 distinct_id。null_as_profile_unset [boolean]默认false。开启后档案属性profile properties中的 null 值会被转换为profile_unset 操作即删除神策档案中该属性已有的值而不是保留旧值。从 SensorsDataSDKWriter.write 与 UserSchemaUtil.buildUnsetUserSchema 的实现看每行记录先执行profileSet随后根据全部属性集合比对出 null 属性并构造 unset schema若所有字段均非 null 则不发送 unset。common optionsSink 插件的通用参数如parallelism、result_table_name等详见 Sink Common Options。四、运行机制从 SeaTunnelRow 到神策 SDK 的完整链路理解数据链路有助于排障。SensorsData Sink 的写入流程可概括为工厂装配SensorsDataSDKSinkFactory 通过AutoService(Factory.class)注册factoryIdentifier()返回SensorsData配置经ReadonlyConfig封装为 SensorsDataSDKSinkConfigSink 创建SensorsDataSDKSink 持有配置与 CatalogTablecreateWriter产出SensorsDataSDKWriterWriter 初始化SensorsDataSDKWriter按consumer创建ConsoleConsumer或BatchConsumer驱动的SensorsAnalytics实例并构建RowAccessor建立列名→索引映射、解析动态事件名、校验列存在性与SensorsDataRecordBuilder逐行转换write(SeaTunnelRow)中由SensorsDataRecordBuilder.Builder.build(row)依据entity_name×record_type构建对应的神策 Schema 对象UserSchema/UserEventSchema/DetailSchema/ item 记录再调用sa.profileSet/sa.track/sa.detailSet/sa.itemSet批量发送SDK 内部按bulk_size/max_cache_row_size阈值批量上报同时prepareCommit()即每次 checkpoint强制sa.flush()。该链路在模块内的测试用例中有完整验证例如 SensorsDataUserRecordTest、SensorsDataSpecialItemRecordTest、TypeUtilTest 与 SensorsDataSDKFactoryTest可以作为理解字段映射与类型转换行为的参考。五、记录类型配置速查根据文档的 Record type requirements四种形态的必填配置组合如下场景entity_namerecord_type必填参数用户事件userseventsevent_name、time_column、distinct_id_column、identity_fields、property_fields用户档案usersusersschema即 users 表、distinct_id_column、identity_fields、property_fields用户明细usersdetailsdetail_id_column、distinct_id_column、identity_fields、property_fields物品记录itemsitemsitem_id_column、item_type_column、property_fields六、完整配置示例以下示例均来自官方文档可直接复制到 SeaTunnel 作业的sink段使用。示例中的server_url为内网演示地址请替换为你自己的神策接收地址。6.1 基础事件埋点Basic Event Trackingsink { SensorsData { server_url http://10.1.136.63:8106/sa?projectdefault time_free true record_type events schema events event_name $AppStart time_column col_date distinct_id_column col_id identity_fields [ { source col_id, target $identity_login_id } { source col_id, target $identity_distinct_id } ] property_fields [ { target prop1, source col1, type INT } { target prop2, source col2, type BIGINT } { target prop3, source col3, type STRING } { target prop4, source col4, type BOOLEAN } ] skip_error_record true } }6.2 动态事件名Dynamic Event Namessink { SensorsData { server_url http://10.1.136.63:8106/sa?projectdefault time_free true record_type events schema events event_name ${event_type} # 事件名取自每行数据 time_column event_timestamp distinct_id_column user_id identity_fields [ { source user_id, target $identity_login_id } { source user_id, target $identity_distinct_id } ] property_fields [ { target price, source amount, type DECIMAL } { target category, source product_category, type STRING } { target device, source device_type, type STRING } ] instant_events [$AppStart, $AppEnd] # 将指定事件标记为即时事件 } }6.3 用户档案更新Profile Property Updates用户档案记录record_type users会更新神策用户档案上挂载的属性。设置null_as_profile_unset true后null 属性会删除对应档案属性而不是保留旧值sink { SensorsData { server_url http://10.1.136.63:8106/sa?projectdefault time_free true entity_name users record_type users schema users distinct_id_column user_id identity_fields [ { source email, target $identity_email } { source phone, target $identity_phone } ] property_fields [ { target name, source full_name, type STRING } { target age, source user_age, type INT } { target gender, source user_gender, type STRING } { target location, source user_location, type STRING } ] null_as_profile_unset true # 属性为 null 时删除该档案属性 } }6.4 物品记录Item Trackingsink { SensorsData { server_url http://10.1.136.63:8106/sa?projectdefault time_free true record_type items schema items event_name $ItemViewed time_column view_time distinct_id_column user_id identity_fields [ { source user_id, target $identity_login_id } ] property_fields [ { target view_duration, source duration, type INT } { target referrer, source referrer_url, type STRING } ] item_id_column product_id item_type_column product_type } }6.5 用户明细记录User Details用户明细记录record_type details为神策用户附加一条带独立标识的明细数据。detail_id_column提供明细键distinct_id_column与identity_fields定位父级用户sink { SensorsData { consumer console server_url http://10.129.27.43:8106/sa?projectsditest time_free true record_type details schema fund_manager distinct_id_column c_id detail_id_column c_id identity_fields [ { target $identity_distinct_id, source c_id } ] property_fields [ { target c_id, source c_id, type STRING } { target fund_amount, source c_int, type INT } { target $is_valid, source c_boolean, type BOOLEAN } ] } }6.6 动态事件名 多身份映射事件名由数据行本身决定event_name ${c_event}并一次性映射多个身份$identity_login_id、$identity_distinct_id。转换失败的记录通过skip_error_record true跳过sink { SensorsData { consumer console server_url http://10.1.136.63:8106/sa?projectdefault time_free true record_type events schema events event_name ${c_event} time_column c_date distinct_id_column c_bigint identity_fields [ { source c_bigint, target $identity_login_id } { source c_bigint, target $identity_distinct_id } ] property_fields [ { target c_tinyint, source c_tinyint, type INT } { target c_bigint, source c_bigint, type BIGINT } { target c_int, source c_int, type INT } { target c_boolean, source c_boolean, type BOOLEAN } ] skip_error_record true } }6.7 控制台输出调试Console Output不发送到神策服务端仅将转换后的记录打印到控制台适合验证字段映射与类型转换是否符合预期sink { SensorsData { server_url http://10.1.136.63:8106/sa?projectdefault consumer console # 打印到控制台而非发送到服务端 record_type events schema events event_name $TestEvent time_column timestamp distinct_id_column test_id property_fields [ { target test, source test_field, type STRING } ] } }七、常见问题与排障建议作业启动即报 Field [xxx] not found in source columndistinct_id_column、time_column、property_fields[].source、identity_fields[].source、动态${event_name}中引用的列名与上游 schema 不一致。这是 RowAccessor 初始化期的强校验请核对上游字段名。Unsupported record type / Unsupported entity nameentity_name与record_type的组合不在支持矩阵内如entity_name items搭配record_type users请参照上文的组合表修正。schema / detailId / itemId / itemType is required对应记录的必填参数缺失或值为空按记录类型配置速查补齐。不想让脏数据中断任务设置skip_error_record true异常行会打印 error 日志后跳过同时可结合consumer console在本地观察转换结果。时间字段解析失败TIMESTAMP/DATE 的字符串解析支持yyyy-MM-dd HH:mm:ss.SSS、yyyy-MM-dd HH:mm:ss、yyyy-MM-dd HH:mm、yyyy-MM-dd、yyyyMMdd_HHmmss、yyyyMMdd六种格式超出范围请在上游先用 Transform 规整格式。事件时间语义需要处理时间而非数据自带时间时将time_column设为current_time()需要历史数据无时区语义时开启time_free true。八、版本与变更记录SensorsData Connector 首次随 SeaTunnel2.3.12版本发布PR #9432完整变更记录见 connector-sensorsdata changelog。由于该 Connector 较新使用前建议确认你的 SeaTunnel 发行版本不低于 2.3.12并在plugin_config或对应的 connector 打包清单中确认connector-sensorsdata已被包含。九、延伸阅读Connector V2 功能概念exactly-once / cdc 等Sink 通用选项SensorsData Connector 源码模块SensorsData 变更记录【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考