首页 文章 精选 留言 我的

精选列表

搜索[数据流],共5172篇文章
优秀的个人博客,低调大师

500万台AP潜在大数据流量入口

500万个无线AP分散覆盖在全国各地,医疗、教育、零售等全场景覆盖、1亿个用户肖像和每天5000TB的海量数据……庞大的数据量撑起了端到端联接的根基,给企业带来更多增值服务机会,也为新华三搭建绿洲平台提供最大支撑。 绿洲平台是国内领先的新IT在线运营平台,也是新华三落地大互联战略的最佳实践。通过新华三绿洲平台,合作伙伴可以尽享云端运维、数据分析、应用共享三大类服务,也可以通过线上运维模式,为企业IT应用落地、运维响应提供便利。 500万台AP重塑入口价值 前些年,说起移动互联网,绝大多数人的认识都停留在慢如蜗牛的WAP站点层面。然而短短几年时间,移动互联网就如“旧时王谢堂前燕,飞入寻常百姓家”一般,在人们生活中迅速扎根。 CNNIC发布的《中国互联网络发展状况统计报告》显示,截至2015年6月,我国网民中使用手机上网的比例为88.9%,其中93%的手机用户利用Wi-Fi接入网络。Wi-Fi已成为移动互联网第一流量入口。 “Wi-Fi是一个非常好的流量入口,竞争非常激烈,互联网企业、商业Wi-Fi运营厂商等都在争夺这个市场。但互联网的经营模式并不是从入口本身赚钱,而是通过大数据分析实现流量变现,从而带来更大价值。”新华三集团副总裁孙德和表示。 在这方面,新华三绿洲平台独具优势,能够将分散在各个场景的数据通过Wi-Fi联接起来,打通线上线下的关系链,然后实现广告营销的个性化定制化服务,进而实现流量变现。 绿洲平台流量入口在于通过AP接入,覆盖大规模用户群。华三通信无线产品部部长白浪介绍,当前全国企业WiFi设备约1500万台,其中华三通信占比超过30%,约500万台无线AP,分布在医疗、教育、零售等20多个场景,能够建立端到端的联接,实现丰富、智能的场景化应用,为未来在数据分析应用、商业广告层面提供更多潜在机会。 覆盖和连接用户中,国内场景营销的先行者光音网络是新华三绿洲平台的首批战略合作伙伴,通过绿洲平台500万台潜在的AP收集上来的海量数据分析消费行为,以及描绘的用户画像,光音网络可进行广告精准推送。光音网络反馈称,比以往传统云平台的SaaS应用,绿洲的服务能够更快切入到场景中,且数据更为精准到位。 只联接不收费 全场景覆盖能力,以及“智能”联接用户,对于各大商家获取用户是极具诱惑力的,因为前者满足用户联网需求,后者则可实现业务增值。如果说以上两点是合伙伴选择新华三绿洲平台的最大理由,那么免费合作政策更是“锦上添花”。 合作过程中,新华三采取“只联接,不收费”的合作模式,那么,新华三绿洲平台为什么免费呢? 孙德和表示,绿洲平台是一个服务延伸平台,它能够降低合作伙伴及新华三自身的人力、销售服务成本。其次,绿洲平台生态圈建设属于刚起步阶段,新华三需要和合作伙伴一起去培养一种习惯和趋势,增加用户“粘性”,不断壮大合作队伍。此外,绿洲平台的最大价值在于价值转换,当越来越多的合作伙伴带着新技术融入合作生态圈后,所有收益都会成倍增长,那时候给客户带来的价值更多,进而也实现新华三绿洲平台构建者的价值重塑。 孙德和强调,新华三将专注于基础架构平台的研发和建设,为互联网公司和WiFi运营公司提供一个坚强的基础架构,帮助合作伙伴将精力投身到客户经营层面,但新华三绝不会介入到运营领域。“我们希望搭建一个基于新IT架构的APP Store,向SaaS服务提供商免费开放” 目前,新华三绿洲平台已经拥有微软、HPE金融服务、友盟+、光音网络、迈外迪、万江龙、海海科技、奥菲传媒、哗哗科技、客如云、云南今日游情、北京国创富盛、河南博加等30多家合作伙伴,这些合作伙伴均可以在新华三绿洲平台应用商店进行交易,无需向平台缴纳任何费用。 以全开放心态构建新生态 只联接,不收费,还能实现业务增值,在为客户提供全方位服务的同时,新华三也在搭建集网络、运营、SaaS应用于一体的WiFi联接新生态。 “绿洲平台的价值是重新定义联接,真正要做到泛联接,人与物、端到云的联接,而不再是散乱分布的各个场景。有了联接和数据后,就可以打通线上线下的关系链,做广告的定制服务等。新华三定位的是做大平台的建设者。”孙德和这样定义绿洲平台。 在孙德和看来,绿洲平台是大互联战略的核心组成部分,能更好地促进新IT成长,为新经济带来新动能。明年,新华三拥有的企业覆盖入口,也将逐渐纳入到“绿洲平台”之下,实现网络全覆盖、盘活大数据资源,开放给相应的公司,创造全新的商业模式。 总结 从“小而美”到“大而美”,新华三在大安全、大互联、云计算、大数据等技术及应用领域,不断实践着。此次,新华三绿洲平台的发布,让合作伙伴享受到领先的IT技术,也给万物互联带来了超乎想象的精彩。新华三将以开放的心态,欢迎业界各类合作伙伴一起共建WiFi新生态。 本文出处:畅享网 本文来自云栖社区合作伙伴畅享网,了解相关信息可以关注vsharing.com网站。

优秀的个人博客,低调大师

【架构实战】实时计算架构:Flink处理亿级数据流

## 一、离线报表晚一天,决策就晚一天 2020年,我们的数据报表是T+1: - 每天凌晨跑批 - 早上8点出报表 - 运营看昨天的数据 **问题:** - 618大促时,销量激增 - 运营想看实时数据调整策略 - 但只能看昨天的数据 - 错过最佳调整时机 引入实时计算后: - 数据秒级更新 - 实时大屏展示 - 实时预警 - 运营决策效率提升10倍 --- ## 二、实时计算架构 ### 2.1 Lambda架构 ``` ┌─────────────────────────────────────────────────────────────────┐ │ Lambda架构 │ │ │ │ 数据源 → Kafka ─┬→ Flink实时处理 → 实时结果 │ │ │ │ │ └→ HDFS离线存储 → Spark离线处理 → 离线结果 │ │ │ │ 实时结果 ─┬→ 合并 → 服务层 │ │ 离线结果 ─┘ │ │ │ │ 优点:容错性好,离线数据修正实时数据 │ │ 缺点:维护两套系统,复杂度高 │ │ │ └──────────────────────────────────────────────────────────────────┘ ``` ### 2.2 Kappa架构 ``` ┌─────────────────────────────────────────────────────────────────┐ │ Kappa架构 │ │ │ │ 数据源 → Kafka → Flink实时处理 → 结果存储 → 服务层 │ │ │ │ 特点: │ │ - 只维护一套实时处理系统 │ │ - Kafka作为唯一数据源 │ │ - 通过重放历史数据实现重新计算 │ │ │ │ 优点:架构简单,维护成本低 │ │ 缺点:历史数据重放成本高 │ │ │ └──────────────────────────────────────────────────────────────────┘ ``` --- ## 三、Flink核心应用 ### 3.1 实时数据大屏 ```java /** * 实时订单统计 */ public class OrderStatisticsJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. Kafka数据源 Properties props = new Properties(); props.setProperty("bootstrap.servers", "kafka:9092"); props.setProperty("group.id", "order-statistics"); FlinkKafkaConsumer source = new FlinkKafkaConsumer<>( "order-topic", new OrderDeserializer(), props ); DataStream orders = env.addSource(source); // 2. 实时统计 orders // 按时间窗口分组(每5秒一个窗口) .keyBy(Order::getProductId) .window(TumblingProcessingTimeWindows.of( Time.seconds(5))) // 聚合计算 .aggregate(new OrderAggregateFunction()) // 写入Redis .addSink(new RedisSink<>(...)); env.execute("Order Statistics"); } } /** * 订单聚合函数 */ public class OrderAggregateFunction implements AggregateFunction { @Override public OrderAccumulator createAccumulator() { return new OrderAccumulator(); } @Override public OrderAccumulator add(Order order, OrderAccumulator acc) { acc.addAmount(order.getAmount()); acc.addCount(1); return acc; } @Override public OrderResult getResult(OrderAccumulator acc) { return new OrderResult( acc.getProductId(), acc.getCount(), acc.getAmount(), System.currentTimeMillis() ); } } ``` ### 3.2 实时风控 ```java /** * 实时风控监测 */ public class RiskControlJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. 消费用户行为流 DataStream behaviors = env .addSource(new KafkaSource<>("user-behavior")); // 2. 实时规则检测 behaviors .keyBy(UserBehavior::getUserId) .process(new RiskDetectionProcessFunction()) .addSink(new AlertSink()); env.execute("Risk Control"); } } /** * 风险检测函数 */ public class RiskDetectionProcessFunction extends KeyedProcessFunction { // 状态:用户最近行为列表 private ListState behaviorState; @Override public void open(Configuration parameters) { ListStateDescriptor descriptor = new ListStateDescriptor<>( "behaviors", UserBehavior.class); behaviorState = getRuntimeContext().getListState(descriptor); } @Override public void processElement(UserBehavior behavior, Context ctx, Collector out) throws Exception { // 1. 添加到状态 behaviorState.add(behavior); // 2. 获取最近1分钟的行为 List recentBehaviors = new ArrayList<>(); for (UserBehavior b : behaviorState.get()) { if (System.currentTimeMillis() - b.getTimestamp() b.getType() == BehaviorType.LOGIN_FAIL) .count(); if (loginFailCount > 5) { out.collect(new RiskAlert( behavior.getUserId(), "LOGIN_FAIL_TOO_MANY", loginFailCount)); } // 规则2:1分钟内下单超过10笔 long orderCount = recentBehaviors.stream() .filter(b -> b.getType() == BehaviorType.ORDER) .count(); if (orderCount > 10) { out.collect(new RiskAlert( behavior.getUserId(), "ORDER_TOO_FREQUENT", orderCount)); } // 4. 注册定时器清理过期状态 ctx.timerService().registerProcessingTimeTimer( System.currentTimeMillis() + 60000); } } ``` ### 3.3 实时ETL ```java /** * 实时数据同步 */ public class RealtimeETLJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. Canal数据变更流 DataStream canalStream = env .addSource(new CanalSource("order-table")); // 2. 数据清洗转换 DataStream orders = canalStream .filter(entry -> entry.getEventType() == INSERT || entry.getEventType() == UPDATE) .map(entry -> convertToOrder(entry)); // 3. 写入多个目标 // 3.1 写入ES(搜索) orders.addSink(new ElasticsearchSink<>(...)); // 3.2 写入Redis(缓存) orders.addSink(new RedisSink<>(...)); // 3.3 写入ClickHouse(分析) orders.addSink(new ClickHouseSink<>(...)); env.execute("Realtime ETL"); } } ``` --- ## 四、状态管理 ### 4.1 状态后端配置 ```java // 配置RocksDB状态后端 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 1. 使用RocksDB(大状态) env.setStateBackend(new EmbeddedRocksDBStateBackend()); // 2. 配置Checkpoint env.enableCheckpointing(60000); // 每60秒checkpoint一次 env.getCheckpointConfig().setCheckpointingMode( CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); env.getCheckpointConfig().enableExternalizedCheckpoints( ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION); // 3. 配置Savepoint路径 env.getCheckpointConfig().setCheckpointStorage( "hdfs:///flink/checkpoints"); ``` ### 4.2 状态恢复 ```java // 从Savepoint恢复 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStateBackend(new EmbeddedRocksDBStateBackend()); // 指定Savepoint路径 String savepointPath = "hdfs:///flink/savepoints/savepoint-123"; env.getCheckpointConfig().setCheckpointStorage(savepointPath); // 从Savepoint恢复执行 env.execute("My Job"); ``` --- ## 五、踩坑实录 ### 坑1:背压问题 > **问题**:下游消费慢,导致上游数据积压。 **解决方案**: ```java // 1. 增加并行度 env.setParallelism(8); // 2. 异步Sink orders.addSink(new AsyncRedisSink<>()); // 3. 调整缓冲区大小 env.getConfig().setBufferTimeout(100); ``` ### 坑2:数据倾斜 > **问题**:某些key数据量大,导致热点。 **解决方案**: ```java // 1. 加盐打散 orders.keyBy(order -> order.getUserId() + "_" + RandomUtils.nextInt(0, 10)) .window(...) // 2. 两阶段聚合 orders.keyBy(order -> order.getUserId() + "_" + salt) .window(...) .aggregate(...) // 第一阶段:局部聚合 .keyBy(...) .window(...) .aggregate(...); // 第二阶段:全局聚合 ``` ### 坑3:状态过大 > **问题**:状态无限增长,内存溢出。 **解决方案**: ```java // 1. 使用TTL清理过期状态 StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.hours(24)) .setUpdateType(UpdateType.OnCreateAndWrite) .setStateVisibility(StateVisibility.NeverReturnExpired) .build(); ValueStateDescriptor descriptor = new ValueStateDescriptor<>("state", String.class); descriptor.enableTimeToLive(ttlConfig); // 2. 使用RocksDB状态后端 env.setStateBackend(new EmbeddedRocksDBStateBackend()); ``` ### 坑4:Exactly-Once不生效 > **问题**:重启后数据重复或丢失。 **解决方案**: ```java // 1. 确保Sink支持事务 FlinkKafkaProducer sink = new FlinkKafkaProducer<>( "topic", new SimpleStringSchema(), props, FlinkKafkaProducer.Semantic.EXACTLY_ONCE // 开启事务 ); // 2. 配置Checkpoint env.enableCheckpointing(60000); env.getCheckpointConfig().setCheckpointingMode( CheckpointingMode.EXACTLY_ONCE); ``` --- ## 六、最佳实践 ### 6.1 Flink调优清单 ``` Flink调优清单: □ 并行度 □ 与Kafka分区数匹配 □ 根据数据量调整 □ 内存 □ TaskManager内存充足 □ Network缓冲区足够 □ 状态 □ 大状态用RocksDB □ 配置状态TTL □ Checkpoint □ 间隔合理(1-5分钟) □ 超时时间充足 □ 保留Checkpoint □ 背压 □ 监控背压指标 □ 及时扩容 ``` --- ## 七、总结 实时计算核心要点: | 组件 | 作用 | 选型 | |------|------|------| | 消息队列 | 数据缓冲 | Kafka | | 计算引擎 | 实时处理 | Flink | | 状态存储 | 状态管理 | RocksDB | | 结果存储 | 结果输出 | Redis/ES | **血的教训:** > 实时计算是数据的"高速公路",搭建好了,数据才能流得快。 --- *个人观点,仅供参考*

优秀的个人博客,低调大师

Apache Beam 2.28.0 发布,大数据流处理与批处理编程范式

Apache Beam 2.28.0 已发布,Beam 是一个用于定义和执行数据处理管道的统一编程模型,包括 ETL、批处理与流处理。Beam 项目重点在于数据处理的编程范式和接口定义,并不涉及具体执行引擎的实现,理想情况是基于 Beam 开发的数据处理程序可以执行在任意的分布式计算引擎上。 更新亮点 与 Parquet 支持相关的大量改进(BEAM-11460,BEAM-8202和BEAM-11526) BeamSQL 中的哈希函数 (BEAM-10074) ZetaSQL 中的哈希函数 (BEAM-11624) 使用 HLL Impl 创建 ApproximateDistinct (BEAM-10324) I/Os SpannerIO 支持面向 Numeric 字段使用 BigDecimal (BEAM-11643) 将 Beam schema 支持添加到ParquetIO (BEAM-11526) 支持 ParquetTable Writer (BEAM-8202) GCP BigQuery sink (streaming inserts) 使用 runner 已确定的分片 (BEAM-11408) PubSub 支持类型:TIMESTAMP, DATE, TIME, DATETIME (BEAM-11533) 新特性/改进 ParquetIO 添加readGenericRecords和readFilesGenericRecords方法可以读取具有未知 schema 的文件。详情查看PR-13554和 (BEAM-11460) 添加对 KafkaTableProvider 中thrift 的支持 (BEAM-11482) 添加对 HadoopFormatIO 的支持以跳过key/value 克隆 (BEAM-11457) 在 Convert.to 转换中支持转换为 GenericRecords (BEAM-11571) 支持读取未知 schema 的 Parquet 文件(BEAM-11460) 发布公告

优秀的个人博客,低调大师

Apache Beam 2.27.0 发布,大数据流处理与批处理编程范式

Apache Beam 2.27.0 发布了。Beam 是一个用于定义和执行数据处理管道的统一编程模型,包括 ETL、批处理与流处理。Beam 项目重点在于数据处理的编程范式和接口定义,并不涉及具体执行引擎的实现,理想情况是基于 Beam 开发的数据处理程序可以执行在任意的分布式计算引擎上。 此版本主要更新内容如下: Highlights Java 11 Containers现已随所有 Beam 版本一起发布。 有一个新的转换ReadAllFromBigQuery,可以在管道运行时接收多个从 BigQuery 读取数据的请求。参见PR 13170和BEAM-9650。 I/Os ReadFromMongoDB 现在可以与 MongoDB Atlas (Python) 一起使用(BEAM-11266) ReadFromMongoDB/WriteToMongoDB将屏蔽 display_data (Python) 中的密码(BEAM-11444) 有一个新的转换ReadAllFromBigQuery,可以在管道运行时接收多个从 BigQuery 读取数据的请求。参见PR 13170和BEAM-9650。 New Features / Improvements 依赖于 Hadoop 的 Beam 模块现在已经测试了与 Hadoop 3 的兼容性(BEAM-8569)。(Hive/HCatalog pending) 作为 Apache Beam 发行过程的一部分,现在支持发布 Java 11 SDK 容器镜像。(BEAM-8106) 向 Beam SQL 添加了 Cloud Bigtable Provider 扩展(BEAM-11173,BEAM-11373) 为 thrift data添加了一个 schema provider(BEAM-11338) 在 Dataflow runner 中添加了 combiner packing pipeline 优化(BEAM-10641) Breaking Changes HBaseIO hbase-shaded-client 依赖项现在应该由用户提供(BEAM-9278) amazon-web-services2 中的--region标志已被-awsRegion替代(BEAM-11331) 更新说明:https://github.com/apache/beam/releases/tag/v2.27.0

优秀的个人博客,低调大师

ubuntu16.04安装Storm数据流实时处理系统 集群

root@master:~# wgethttp://mirror.bit.edu.cn/apache/storm/apache-storm-1.1.1/apache-storm-1.1.1.tar.gz root@master:/usr/local/apache-storm-1.1.1# vim conf/storm.yaml storm.zookeeper.servers: - "master" - "slave1" - "slave2" 复制到其它节点 root@master:/usr/local/apache-storm-1.1.1# bin/storm nimbus & [1] 33251 root@slave1:/usr/local/apache-storm-1.1.1# bin/storm supervisor & [1] 15896 root@master:/usr/local/apache-storm-1.1.1# bin/storm ui & [2] 33436 root@slave1:/usr/local/apache-storm-1.1.1# bin/storm ui & [2] 16009 root@master:/usr/local/apache-storm-1.1.1# jps 14033 Master 33251 nimbus 33525 Jps 14901 QuorumPeerMain 13033 SecondaryNameNode 33436 core 12813 NameNode 13182 ResourceManager root@slave1:/usr/local/apache-storm-1.1.1# root@slave1:/usr/local/apache-storm-1.1.1# jps 16096 Jps 15457 HRegionServer 7219 DataNode 15896 Supervisor 8632 QuorumPeerMain 16009 core 8330 Worker 8268 Worker 本文转自 OpenStack2015 博客,原文链接:http://blog.51cto.com/andyliu/1967308 如需转载请自行联系原作者

优秀的个人博客,低调大师

苹果和谷歌联手合作,简化 iPhone 与 Android 双向迁移数据流

根据科技媒体 9to5google 的报道,谷歌与苹果共同开发的“换机助手”已初步体现在最新发布的 Android Canary 开发者测试版中,并预计将在即将推出的i OS 26 开发者版本中同步亮相。 长期以来,两大移动生态系统的用户在更换平台时需依赖第三方工具或官方提供的独立应用——如苹果的“转移到iOS”和谷歌的“Switch to Android”。尽管这些应用已能迁移部分数据(如联系人、照片和日历),但在完整性、易用性和兼容性方面仍存在局限。 报道称,经谷歌代表证实,新功能将直接嵌入设备初始设置流程,使用户在激活新手机时即可更高效、安全地从另一平台导入更多类型的数据,包括消息记录、应用设置乃至部分媒体内容。 目前,该功能尚处于早期开发阶段,仅在 Android Canary 版本中可见雏形,具体支持的数据类型、传输协议及隐私保护机制等细节尚未公开。 相关阅读:苹果开发新框架 AppMigrationKit,让 iPhone 与 Android 数据自由迁移

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册